"""Dashboard 统计服务 — 上帝视角(全厂全系统数据,不按用户过滤)""" from datetime import datetime from sqlalchemy import select, func, or_ from sqlalchemy.ext.asyncio import AsyncSession from pydantic import BaseModel # ============================================================ # Schemas # ============================================================ class DashboardStats(BaseModel): products_total: int products_pending: int products_in_progress: int products_completed: int tasks_total: int tasks_pending: int tasks_in_progress: int tasks_completed: int tasks_rejected: int tasks_rework: int unread_notifications: int = 0 unread_messages: int = 0 class WipTask(BaseModel): task_id: str task_name: str assignee: str product_sn: str # 16位HEX身份证 external_serial: str | None # 业务序列号 material_name: str # 设备名称 spec_model: str # 规格型号 status: str received_at: str duration_hours: float class CompletedTask(BaseModel): task_id: str task_name: str # 完成的工序节点 assignee: str # 完成人中文姓名 product_sn: str # 16位HEX身份证 external_serial: str | None # 业务序列号 material_name: str # 产品名称(物料名称) spec_model: str # 规格型号 completed_at: str # 完成时间 ISO class RejectedTask(BaseModel): task_id: str kind: str = "rejected" # "rejected"(已驳回) | "rework"(返工任务) task_name: str # 工序 product_sn: str # 16位HEX身份证 external_serial: str | None # 业务序列号 material_name: str # 产品名称(物料名称) spec_model: str # 规格型号 rejected_by: str = "—" # 驳回人中文姓名(rework 可能无) rework_assignee: str # 返工任务负责人中文姓名(转交给谁返工) reject_reason: str | None # 驳回原因 status: str | None = None # 当前状态(rework 任务用:PENDING/WIP/COMPLETED) rejected_at: str | None # 驳回时间 ISO(BEIJING_TZ) class UserOperation(BaseModel): user_id: str # 登录名 username user_name: str # 中文姓名 receive_count: int = 0 # 接收次数 transfer_count: int = 0 # 转交次数(action=complete) record_count: int = 0 # 上传备注次数(action=record) total: int = 0 # 总操作次数 class OperationDetail(BaseModel): task_name: str # 任务/工序名 product_sn: str # 产品身份证 material_name: str # 设备名称 remark: str | None # 备注/说明 time: str # 操作时间 ISO(BEIJING_TZ) class WipMatrixRow(BaseModel): spec_model: str # 规格型号(Y 轴) product_name: str = "" # 产品名称(物料名称,用于首列复合显示) dimension_key: str # 人员姓名 或 工序名称(X 轴) count: int = 0 # 该交叉点的设备数量 assignees: list[str] = [] # 该交叉点涉及的主负责人(中文名,去重) class WipMatrixDetailRow(BaseModel): """WIP 矩阵单元格下钻 — 某规格型号 × 某工序下的设备明细""" product_id: str serial_number: str # 16位HEX身份证 external_serial: str | None = None # 业务序列号 material_name: str = "" # 产品名称 spec_model: str = "" task_status: str = "" # 当前状态 WIP/PENDING/COMPLETED/ARCHIVED/OUTBOUND assignee_id: str = "" assignee: str = "" # 负责人中文名 duration_hours: float = 0.0 # 已在该工序滞留时长(小时) total_hours: float = 0.0 # 整个项目总时长(自然小时,最早介入→最晚结束) work_hours: float = 0.0 # 仅工作日有效时长(小时,排除周末/节假日) received_at: str | None = None # 接手时间 MM-DD HH:mm class PersonDevice(BaseModel): product_id: str serial_number: str # 16位HEX身份证 external_serial: str | None # 业务序列号 material_name: str # 产品名称 spec_model: str # 规格型号 task_status: str # 该设备名下的状态 WIP/PENDING duration_hours: float # 滞留时长(在当前人手上多久,小时) received_at: str | None # 接收时间(北京时间 MM-DD HH:mm) class PersonWorkload(BaseModel): assignee_id: str assignee_name: str # 中文姓名 device_count: int devices: list[PersonDevice] class PersonHistoryRecord(BaseModel): task_id: str task_name: str # 工序/任务名 assignee_id: str # 负责人ID assignee_name: str # 负责人姓名 product_sn: str # 16位身份证 external_serial: str | None # 业务序列号 material_name: str # 产品名称 spec_model: str # 规格型号 status: str # 状态 WIP/PENDING/COMPLETED received_at: str | None # 开始时间 MM-DD HH:mm completed_at: str | None # 结束时间 MM-DD HH:mm(进行中为 None) duration_hours: float # 总耗时(小时) latest_valid_remark: str | None = None # 最新有效备注(排除系统转交类) record_count: int = 0 # 记录总数(含系统备注) class ProductMessageItem(BaseModel): id: str content: str operator_name: str # 留言人中文姓名 product_sn: str # 16位HEX系统追溯码 external_serial: str | None # 业务产品序列号(如 25022) material_name: str # 物料名称 created_at: str # ISO时间字符串 class ProductMessageList(BaseModel): items: list[ProductMessageItem] total: int # ============================================================ # 看板统计(时间快照语义) # ============================================================ async def get_dashboard_stats( db: AsyncSession, since: datetime | None = None, until: datetime | None = None, ) -> DashboardStats: """ 上帝视角 — 全厂全系统统计。 时间筛选规则: - PENDING / WIP / 总数:永远忽略时间筛选,返回实时快照。 - COMPLETED / REJECTED / 完成率:严格按 since~until 过滤(用于时段报表)。 """ from app.models.product import Product from app.models.task import ( Task, TASK_STATUS_PENDING, TASK_STATUS_WIP, TASK_STATUS_COMPLETED, TASK_STATUS_REJECTED, ) from app.models.notification import Notification from app.models.message import ProductMessage # ── 产品(实时快照,不过滤) ── # 状态口径(用户确认):废弃 Product.status 恒值判断, # COMPLETED=待仓库收货,ARCHIVED=已入库('在库' 旧命名等价),OUTBOUND=已出库。 # 产品流转卡片第三段"已完结"= 待收货 + 已实收 + 旧在库 + 已出库,保证三段和 = 总数。 p_total = await db.scalar(select(func.count(Product.id))) finished_cond = Product.overall_status.in_(["待仓库收货", "已入库", "在库", "已出库"]) p_done = await db.scalar(select(func.count(Product.id)).where(finished_cond)) # 在制 WIP = 未完结 且 存在活跃任务(PENDING/WIP) 的产品数 not_finished = or_(Product.overall_status.is_(None), ~finished_cond) has_active_task = select(Task.id).where( Task.product_id == Product.id, Task.status.in_([TASK_STATUS_PENDING, TASK_STATUS_WIP]), ).exists() p_progress = await db.scalar( select(func.count(Product.id)).where(not_finished, has_active_task) ) # 待流转 = 总数 - 在制 - 完结(三段互斥,保证进度条总和=总数) p_pending = max((p_total or 0) - (p_progress or 0) - (p_done or 0), 0) # ── 任务实时快照(PENDING/WIP/返工 — 永远不过滤) ── t_pending = await db.scalar(select(func.count(Task.id)).where(Task.status == TASK_STATUS_PENDING)) t_progress = await db.scalar(select(func.count(Task.id)).where(Task.status == TASK_STATUS_WIP)) t_rework = await db.scalar(select(func.count(Task.id)).where(Task.is_rework.is_(True))) # ── 任务已完成/驳回(时间可过滤) ── t_done_q = select(func.count(Task.id)).where(Task.status == TASK_STATUS_COMPLETED) t_rej_q = select(func.count(Task.id)).where(Task.status == TASK_STATUS_REJECTED) if since: t_done_q = t_done_q.where(Task.completed_at >= since) t_rej_q = t_rej_q.where(Task.completed_at >= since) if until: t_done_q = t_done_q.where(Task.completed_at <= until) t_rej_q = t_rej_q.where(Task.completed_at <= until) t_done = await db.scalar(t_done_q) t_rejected = await db.scalar(t_rej_q) # ── 任务总数 = 实时快照段 + 时间过滤段(确保进度条段总和=总数) ── t_total = (t_pending or 0) + (t_progress or 0) + (t_done or 0) + (t_rejected or 0) # ── 通知 & 留言(实时快照) ── unread_notif = await db.scalar( select(func.count(Notification.id)).where(Notification.is_read.is_(False)) ) unread_msg = await db.scalar(select(func.count(ProductMessage.id))) return DashboardStats( products_total=p_total or 0, products_pending=p_pending or 0, products_in_progress=p_progress or 0, products_completed=p_done or 0, tasks_total=t_total or 0, tasks_pending=t_pending or 0, tasks_in_progress=t_progress or 0, tasks_completed=t_done or 0, tasks_rejected=t_rejected or 0, tasks_rework=t_rework or 0, unread_notifications=unread_notif or 0, unread_messages=unread_msg or 0, ) # ============================================================ # 在制品看板(永远实时) # ============================================================ async def get_wip_tasks(db: AsyncSession, limit: int = 20) -> list[WipTask]: from app.models.task import Task, TASK_STATUS_PENDING, TASK_STATUS_WIP from app.models.product import Product from app.models.holiday import Holiday from app.core.time_utils import get_beijing_time, to_beijing, working_duration_hours # 读取节假日(排除非工作日) hres = await db.execute(select(Holiday.day)) holidays = {r[0] for r in hres} stmt = ( select(Task, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model) .join(Product, Task.product_id == Product.id) .where(Task.status.in_([TASK_STATUS_PENDING, TASK_STATUS_WIP])) .limit(limit * 2) ) result = await db.execute(stmt) rows = result.all() raw_ids = list({t.assignee_id for t, *_ in rows if t.assignee_id}) name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(raw_ids) now = get_beijing_time() wip_list: list[WipTask] = [] for task, product_sn, ext_sn, mat_name, spec in rows: start = task.received_at or task.created_at if start: start_bj = to_beijing(start) # 🚀 naive 按 UTC 转北京时间(修复多算8小时) hours = working_duration_hours(start_bj, now, holidays) recv_str = start_bj.strftime("%m-%d %H:%M") else: hours = 0 recv_str = "" wip_list.append(WipTask( task_id=str(task.id), task_name=task.task_name, assignee=name_map.get(task.assignee_id or "", task.assignee_id or "未分配"), product_sn=product_sn or "", external_serial=ext_sn or None, material_name=mat_name or "", spec_model=spec or "", status=task.status, received_at=recv_str, duration_hours=hours, )) # 统一按滞留时间降序排列(无视 status,纯数值排序) wip_list.sort(key=lambda t: t.duration_hours, reverse=True) return wip_list[:limit] # ============================================================ # 已完成任务明细(流转完成率下钻 — 按时段过滤) # ============================================================ async def get_completed_tasks( db: AsyncSession, since: datetime | None = None, until: datetime | None = None, limit: int = 200, ) -> list[CompletedTask]: """按时段查询已完成任务明细(上帝视角),用于「流转完成率」卡片下钻。""" from app.models.task import Task, TASK_STATUS_COMPLETED from app.models.product import Product from app.core.time_utils import BEIJING_TZ stmt = ( select(Task, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model) .join(Product, Task.product_id == Product.id) .where(Task.status == TASK_STATUS_COMPLETED) ) if since: stmt = stmt.where(Task.completed_at >= since) if until: stmt = stmt.where(Task.completed_at <= until) stmt = stmt.order_by(Task.completed_at.desc()).limit(limit) result = await db.execute(stmt) rows = result.all() raw_ids = list({t.assignee_id for t, *_ in rows if t.assignee_id}) name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(raw_ids) items: list[CompletedTask] = [] for task, sn, ext, mat, spec in rows: t = task.completed_at if t: if t.tzinfo is None: from datetime import timezone as dt_timezone t = t.replace(tzinfo=dt_timezone.utc).astimezone(BEIJING_TZ) else: t = t.astimezone(BEIJING_TZ) time_str = t.isoformat() if t else "" items.append(CompletedTask( task_id=str(task.id), task_name=task.task_name, assignee=name_map.get(task.assignee_id or "", task.assignee_id or "未分配"), product_sn=sn or "", external_serial=ext or None, material_name=mat or "", spec_model=spec or "", completed_at=time_str, )) return items # ============================================================ # 被驳回任务明细(品质驳回下钻 — 按时段过滤) # ============================================================ async def get_rejected_tasks( db: AsyncSession, since: datetime | None = None, until: datetime | None = None, limit: int = 200, ) -> list[RejectedTask]: """查询驳回/返工任务明细(上帝视角),用于「驳回/返工」卡片下钻。 返回两类(与卡片数字 tasks_rejected + tasks_rework 口径一致): - kind="rejected":被驳回任务,按 completed_at 时间过滤 - kind="rework":返工任务(is_rework=True 且非驳回状态),实时快照不过滤时间 返工负责人追溯逻辑与 task_service.reject_task 一致: 优先最早 create log 的 operator → 兜底父任务负责人 → 兜底自身。 """ from app.models.task import Task, TASK_STATUS_REJECTED from app.models.task_log import TaskLog from app.models.product import Product from app.core.time_utils import BEIJING_TZ def _to_bj_iso(dt) -> str: if not dt: return "" if dt.tzinfo is None: from datetime import timezone as dt_timezone dt = dt.replace(tzinfo=dt_timezone.utc).astimezone(BEIJING_TZ) else: dt = dt.astimezone(BEIJING_TZ) return dt.isoformat() # ── 1. 已驳回任务(按时间过滤)── stmt = ( select(Task, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model) .join(Product, Task.product_id == Product.id) .where(Task.status == TASK_STATUS_REJECTED) ) if since: stmt = stmt.where(Task.completed_at >= since) if until: stmt = stmt.where(Task.completed_at <= until) stmt = stmt.order_by(Task.completed_at.desc()).limit(limit) rejected_rows = (await db.execute(stmt)).all() # ── 2. 返工任务(is_rework=True 且当前非驳回状态,实时不过滤时间)── rework_stmt = ( select(Task, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model) .join(Product, Task.product_id == Product.id) .where(Task.is_rework.is_(True), Task.status != TASK_STATUS_REJECTED) .order_by(Task.created_at.desc()) .limit(limit) ) rework_rows = (await db.execute(rework_stmt)).all() # 合并:rejected 在前,rework 在后 all_rows: list[tuple[str, object]] = [ *[("rejected", t) for t in rejected_rows], *[("rework", t) for t in rework_rows], ] # ── 已驳回任务:批量取驳回人 / 最早 create log / 父任务负责人 ── rejected_ids = [t.id for t, *_ in rejected_rows] reject_op: dict = {} if rejected_ids: sub = ( select( TaskLog.task_id, TaskLog.operator_id, func.row_number().over( partition_by=TaskLog.task_id, order_by=TaskLog.created_at.desc(), ).label("rn"), ) .where(TaskLog.task_id.in_(rejected_ids), TaskLog.action_type == "reject") ).subquery() r = await db.execute(select(sub.c.task_id, sub.c.operator_id).where(sub.c.rn == 1)) for row in r: if row[1]: reject_op[row[0]] = row[1] create_op: dict = {} if rejected_ids: sub = ( select( TaskLog.task_id, TaskLog.operator_id, func.row_number().over( partition_by=TaskLog.task_id, order_by=TaskLog.created_at.asc(), ).label("rn"), ) .where(TaskLog.task_id.in_(rejected_ids), TaskLog.action_type == "create") ).subquery() r = await db.execute(select(sub.c.task_id, sub.c.operator_id).where(sub.c.rn == 1)) for row in r: if row[1]: create_op[row[0]] = row[1] parent_assignee: dict = {} parent_ids = [t.parent_task_id for t, *_ in rejected_rows if t.parent_task_id and t.id not in create_op] if parent_ids: r = await db.execute( select(Task.id, Task.assignee_id).where(Task.id.in_(parent_ids)) ) for row in r: if row[1]: parent_assignee[row[0]] = row[1] # ── 中文名映射(驳回人 + 返工负责人 一次批量查)── raw_ids: set[str] = set() for kind, (t, *_row) in all_rows: if kind == "rejected": raw_ids.add(reject_op.get(t.id) or "") raw_ids.add(create_op.get(t.id) or "") if t.id not in create_op and t.parent_task_id: raw_ids.add(parent_assignee.get(t.parent_task_id) or "") raw_ids.add(t.assignee_id or "") raw_ids.discard("") name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(list(raw_ids)) items: list[RejectedTask] = [] for kind, (task, sn, ext, mat, spec) in all_rows: if kind == "rejected": # 复刻 reject_task 追溯逻辑:create op → 父任务负责人 → 自身 rework_id = create_op.get(task.id) if not rework_id and task.parent_task_id: rework_id = parent_assignee.get(task.parent_task_id) if not rework_id: rework_id = task.assignee_id rejected_by_id = reject_op.get(task.id) or task.assignee_id items.append(RejectedTask( task_id=str(task.id), kind="rejected", task_name=task.task_name, product_sn=sn or "", external_serial=ext or None, material_name=mat or "", spec_model=spec or "", rejected_by=name_map.get(rejected_by_id or "", rejected_by_id or "—"), rework_assignee=name_map.get(rework_id or "", rework_id or "—"), reject_reason=task.reject_reason, status=task.status, rejected_at=_to_bj_iso(task.completed_at), )) else: # 返工任务:负责人即其 assignee items.append(RejectedTask( task_id=str(task.id), kind="rework", task_name=task.task_name, product_sn=sn or "", external_serial=ext or None, material_name=mat or "", spec_model=spec or "", rejected_by="—", rework_assignee=name_map.get(task.assignee_id or "", task.assignee_id or "—"), reject_reason=None, status=task.status, rejected_at=_to_bj_iso(task.received_at or task.created_at), )) return items # ============================================================ # 人员操作统计(接收/转交/上传备注 — 按人聚合,时间可筛选) # ============================================================ async def get_user_operations( db: AsyncSession, since: datetime | None = None, until: datetime | None = None, ) -> list[UserOperation]: """上帝视角 — 统计**全部人员**的操作次数(按时段过滤)。 返回所有 IRIS 部门人员(该时段无操作的计 0),并追加有操作记录但不在 人员清单中的账号(如历史/已离职)。 操作口径(均不修改数据库): - 接收: task_logs.action_type='receive'(按操作人) - 转交: task_logs.action_type='complete'(按操作人) - 上传备注: task_records 按**任务负责人**归因(排除系统自动生成的 以 '[' 开头的备注,如 "[接收] 操作员已确认接收"),历史数据可回溯 """ from app.models.task_log import TaskLog from app.models.task import Task, TaskRecord from app.core.mom_database import MomSessionLocal from sqlalchemy import text, or_ # ── 1. 接收 / 转交(task_logs 按 operator_id 聚合)── rcv = func.count().filter(TaskLog.action_type == "receive") cpl = func.count().filter(TaskLog.action_type == "complete") stmt = ( select(TaskLog.operator_id, rcv.label("receive"), cpl.label("complete")) .where( TaskLog.action_type.in_(["receive", "complete"]), TaskLog.operator_id.isnot(None), ) .group_by(TaskLog.operator_id) ) if since: stmt = stmt.where(TaskLog.created_at >= since) if until: stmt = stmt.where(TaskLog.created_at <= until) op_rows = (await db.execute(stmt)).all() op_map = {r[0]: (r[1] or 0, r[2] or 0) for r in op_rows} # ── 2. 上传备注(task_records 按任务 assignee 归因,排除系统自动备注)── rcd_stmt = ( select(Task.assignee_id, func.count(TaskRecord.id)) .join(Task, TaskRecord.task_id == Task.id) .where( Task.assignee_id.isnot(None), or_(TaskRecord.remark.is_(None), ~TaskRecord.remark.like("[%")), ) .group_by(Task.assignee_id) ) if since: rcd_stmt = rcd_stmt.where(TaskRecord.created_at >= since) if until: rcd_stmt = rcd_stmt.where(TaskRecord.created_at <= until) rcd_rows = (await db.execute(rcd_stmt)).all() record_map = {r[0]: r[1] or 0 for r in rcd_rows} # ── 3. 获取全部 MOM 用户清单(IRIS 部门)── users: list[dict] = [] try: dbm = MomSessionLocal() try: rows = dbm.execute(text(""" SELECT username, SPLIT_PART(username, '/', 1) AS full_name FROM sys_user WHERE department = 'IRIS' """)).fetchall() except Exception: rows = dbm.execute(text(""" SELECT username, SPLIT_PART(username, '/', 1) AS full_name FROM sys_user """)).fetchall() finally: dbm.close() for row in rows: short = row.username.split("/")[-1] if "/" in row.username else row.username users.append({"username": short, "full_name": row.full_name or short}) except Exception: users = [] # ── 4. 合并:所有人员 + 有操作但不在清单的账号 ── merged: dict[str, dict] = {} for u in users: merged[u["username"]] = {"name": u["full_name"], "rcv": 0, "cpl": 0, "rcd": 0} for uid in set(op_map.keys()) | set(record_map.keys()): if uid not in merged: merged[uid] = {"name": uid, "rcv": 0, "cpl": 0, "rcd": 0} r, c = op_map.get(uid, (0, 0)) merged[uid]["rcv"] = r merged[uid]["cpl"] = c merged[uid]["rcd"] = record_map.get(uid, 0) # ── 5. 组装 + 排序(总次数降序,无操作的排在最后)── items: list[UserOperation] = [] for uid, v in merged.items(): items.append(UserOperation( user_id=uid, user_name=v["name"], receive_count=v["rcv"], transfer_count=v["cpl"], record_count=v["rcd"], total=v["rcv"] + v["cpl"] + v["rcd"], )) items.sort(key=lambda x: x.total, reverse=True) return items # ============================================================ # 在制品分布透视表(WIP Matrix:规格型号 × 人员/工序) # ============================================================ async def get_wip_matrix( db: AsyncSession, dimension: str = "assignee", since: datetime | None = None, until: datetime | None = None, ) -> list[WipMatrixRow]: """生产分布透视表:Y=规格型号,X=人员 或 工序,单元格=设备数量。 核心口径:**每台设备只统计一次**,按它「当前所处工序」归属—— 1. 取该设备最新的一条主分支任务(parent_task_id IS NULL 或 TRANSFER/RECOVERY) 2. 设备当前工序 = 该最新主任务的工序名(task_name),不区分状态 - 「待确认」(PENDING) = 别人转给我但未接收 - 「已入库」= MOM 已扫码实收;「已完成」= 车间完工待实收 - 其他已完成工序(如「测试」完成)按原工序显示,不强制归仓库态 这样一台设备在「上一步已完成 + 下一步待确认」时只算一次(待确认),不会重复计数。 since/until 时间过滤为“动静分离”:仅对 终结(收口/在库/出库) 设备过滤,锚点取 最新主线任务的 completed_at(完工/变动时间,兜底 created_at,schema 无 updated_at); 在制(未收口)设备无条件全量计入,保证与人员看板实时在制不脱节。 dimension: - assignee: 按当前任务负责人聚合(dimension_key 为中文姓名) - task_name: 按当前工序聚合(dimension_key 为工序名,附主负责人) """ from datetime import timezone as dt_timezone from app.models.task import Task from app.models.product import Product # 每台设备按主任务创建时间倒序,取第一条即「最新主任务」 result = await db.execute( select( Product.id, Product.spec_model, Product.material_name, Task.task_name, Task.assignee_id, Task.created_at, Task.completed_at, Product.current_location_id, Product.overall_status, Product.status, ) .join(Task, Task.product_id == Product.id) .where( or_( Task.parent_task_id.is_(None), Task.task_type.in_(["TRANSFER", "RECOVERY", "WAREHOUSE"]), ) ) .order_by(Product.id, Task.created_at.desc()) ) rows = result.all() # 每台设备 → (spec, 当前工序/负责人, 当前负责人ID) device_cur: dict[str, tuple] = {} seen: set[str] = set() for pid, spec, product_name, task_name, assignee, created, completed, loc, overall_status, product_status in rows: if pid in seen: continue seen.add(pid) # ── 真实产出时间锚点:完工/变动时间优先,创建时间仅兜底 ── # 说明:tasks 表无 updated_at 列(已核实 information_schema),故锚点取 # completed_at or created_at;若后续新增 updated_at 可并入。 anchor = completed or created # ── 动静分离:仅终结(收口/在库/出库)设备按 since~until 过滤;在制设备无条件全量 ── is_terminal = ( loc == "virtual_warehouse" or overall_status in ("待仓库收货", "已入库", "在库", "已出库") or str(product_status).upper() in ("ARCHIVED", "OUTBOUND") or task_name in ("扫码入库", "扫码出库") ) if is_terminal: if anchor is not None and anchor.tzinfo is None: anchor = anchor.replace(tzinfo=dt_timezone.utc) if since and anchor is not None and anchor < since: continue if until and anchor is not None and anchor > until: continue # else: 在制设备无视时间参数,保证与人员看板“实时在制”总数一致 if dimension == "task_name": # ── 终结三列判定(多维防御,不依赖单一 loc / overall)── # 已入库: 整体实收(已入库/在库) 或 status==ARCHIVED 或 收口节点「扫码入库」 if overall_status in ("已入库", "在库") or str(product_status).upper() == "ARCHIVED" or task_name == "扫码入库": key = "已入库" # 已出库: 整体已出库 或 status==OUTBOUND 或 收口节点「扫码出库」 elif overall_status == "已出库" or str(product_status).upper() == "OUTBOUND" or task_name == "扫码出库": key = "已出库" # 已完成: 车间完工待实收(待仓库收货) 或 处于仓库(未实收/未出库) elif overall_status == "待仓库收货" or loc == "virtual_warehouse": key = "已完成" else: key = task_name or "—" else: key = assignee or "未分配" device_cur[pid] = (spec or "未知型号", key, assignee or "", product_name or "") # 聚合:规格 × 当前工序 → 设备数;同时收集负责人 与 产品名称 agg: dict[tuple, int] = {} assignee_map: dict[tuple, set] = {} product_name_map: dict[tuple, str] = {} for spec, key, assignee, product_name in device_cur.values(): k = (spec, key) agg[k] = agg.get(k, 0) + 1 product_name_map.setdefault(k, product_name) if assignee: assignee_map.setdefault(k, set()).add(assignee) # 负责人 ID → 中文名 raw_ids: set[str] = set() for s in assignee_map.values(): raw_ids |= s name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(list(raw_ids)) items: list[WipMatrixRow] = [] for (spec, key), cnt in agg.items(): dim_display = name_map.get(key, key) if dimension == "assignee" else key assignees = [name_map.get(a, a) for a in assignee_map.get((spec, key), set())] or [] items.append(WipMatrixRow( spec_model=spec, product_name=product_name_map.get((spec, key), ""), dimension_key=dim_display, count=cnt, assignees=assignees, )) items.sort(key=lambda x: (x.spec_model, x.dimension_key)) return items # ============================================================ # WIP 矩阵单元格下钻 — 某规格型号 × 某工序下的设备明细 # ============================================================ async def get_wip_matrix_detail( db: AsyncSession, spec_model: str, process: str, since: datetime | None = None, until: datetime | None = None, ) -> list[WipMatrixDetailRow]: """WIP 矩阵单元格下钻:返回该 规格型号×工序 交叉点下的设备明细。 口径与 get_wip_matrix 完全一致:每台设备取最新一条主分支任务, 按其在库三态(已完成/已入库/已出库)或工序名归类到 process。 """ from datetime import timezone as dt_timezone from app.models.task import Task from app.models.product import Product from app.models.holiday import Holiday from app.core.time_utils import get_beijing_time, to_beijing, working_duration_hours hres = await db.execute(select(Holiday.day)) holidays = {r[0] for r in hres} result = await db.execute( select( Product.id, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model, Product.current_location_id, Product.overall_status, Product.status, Task.task_name, Task.assignee_id, Task.status.label("task_status"), Task.created_at, Task.completed_at, Task.received_at, ) .join(Task, Task.product_id == Product.id) .where( or_( Task.parent_task_id.is_(None), Task.task_type.in_(["TRANSFER", "RECOVERY", "WAREHOUSE"]), ) ) .order_by(Product.id, Task.created_at.desc()) ) rows = result.all() now = get_beijing_time() seen: set[str] = set() matched: list[WipMatrixDetailRow] = [] raw_names: set[str] = set() for pid, serial, ext, mat_name, spec, loc, overall, pstatus, task_name, assignee, tstatus, created, completed, received in rows: if pid in seen: continue seen.add(pid) if spec_model and (spec or "") != spec_model: continue # ── 真实产出时间锚点:完工/变动时间优先,创建时间仅兜底(schema 无 updated_at)── anchor = completed or created # ── 动静分离:仅终结(收口/在库/出库)设备按 since~until 过滤;在制设备无条件全量 ── is_terminal = ( loc == "virtual_warehouse" or overall in ("待仓库收货", "已入库", "在库", "已出库") or str(pstatus).upper() in ("ARCHIVED", "OUTBOUND") or task_name in ("扫码入库", "扫码出库") ) if is_terminal: if anchor is not None and anchor.tzinfo is None: anchor = anchor.replace(tzinfo=dt_timezone.utc) if since and anchor is not None and anchor < since: continue if until and anchor is not None and anchor > until: continue # else: 在制设备无视时间参数,保证与人员看板“实时在制”总数一致 # ── 终结三列判定(与 get_wip_matrix 一字不差同步,多维防御)── if overall in ("已入库", "在库") or str(pstatus).upper() == "ARCHIVED" or task_name == "扫码入库": key = "已入库" elif overall == "已出库" or str(pstatus).upper() == "OUTBOUND" or task_name == "扫码出库": key = "已出库" elif overall == "待仓库收货" or loc == "virtual_warehouse": key = "已完成" else: key = task_name or "—" if process and key != process: continue # 滞留时长(当前工序接手时间 → 现在) start = received or created duration = 0.0 received_str = None if start: start_bj = to_beijing(start) duration = working_duration_hours(start_bj, now, holidays) received_str = start_bj.strftime("%m-%d %H:%M") matched.append(WipMatrixDetailRow( product_id=str(pid), serial_number=serial or "", external_serial=ext, material_name=mat_name or "", spec_model=spec or "", task_status=tstatus or "", assignee_id=assignee or "", duration_hours=duration, received_at=received_str, )) if assignee: raw_names.add(assignee) # 负责人中文名映射 name_map: dict[str, str] = {} if raw_names: from app.services.mom_cache import get_display_names name_map = get_display_names(list(raw_names)) for r in matched: r.assignee = name_map.get(r.assignee_id, r.assignee_id) # ── 整个项目总时长:设备最早介入 → 最晚结束(自然 + 工作日) ── if matched: import uuid as uuid_mod from sqlalchemy import func as sa_func uuid_list = [uuid_mod.UUID(m.product_id) for m in matched] range_rows = (await db.execute( select( Task.product_id, sa_func.min(sa_func.coalesce(Task.received_at, Task.created_at)).label("t0"), sa_func.max(sa_func.coalesce(Task.completed_at, now)).label("tmax"), ) .where(Task.product_id.in_(uuid_list)) .group_by(Task.product_id) )).all() row_by_pid = {m.product_id: m for m in matched} for pid, t0, tmax in range_rows: row = row_by_pid.get(str(pid)) if row is None or t0 is None or tmax is None: continue row.total_hours = round((tmax - t0).total_seconds() / 3600, 1) row.work_hours = round(working_duration_hours(to_beijing(t0), to_beijing(tmax), holidays), 1) return matched # ============================================================ # 人员操作明细(点击数字下钻 — 接收/转交/上传备注) # ============================================================ async def get_user_operation_detail( db: AsyncSession, user_id: str, action_type: str, since: datetime | None = None, until: datetime | None = None, ) -> list[OperationDetail]: """查询某人在指定时段内的某类操作明细。 action_type: - receive: 接收(task_logs.action_type='receive') - transfer: 转交(task_logs.action_type='complete') - record: 上传备注(该人名下任务的手动备注,排除系统自动生成的) """ from app.models.task_log import TaskLog from app.models.task import Task, TaskRecord from app.models.product import Product from app.core.time_utils import BEIJING_TZ def _to_iso(dt) -> str: if not dt: return "" if dt.tzinfo is None: from datetime import timezone as dt_timezone dt = dt.replace(tzinfo=dt_timezone.utc).astimezone(BEIJING_TZ) else: dt = dt.astimezone(BEIJING_TZ) return dt.isoformat() items: list[OperationDetail] = [] if action_type in ("receive", "transfer"): act = "receive" if action_type == "receive" else "complete" stmt = ( select(TaskLog.created_at, Task.task_name, Product.serial_number, Product.material_name, TaskLog.remark) .join(Task, TaskLog.task_id == Task.id) .join(Product, Task.product_id == Product.id) .where(TaskLog.operator_id == user_id, TaskLog.action_type == act) .order_by(TaskLog.created_at.desc()) ) if since: stmt = stmt.where(TaskLog.created_at >= since) if until: stmt = stmt.where(TaskLog.created_at <= until) rows = (await db.execute(stmt)).all() for r in rows: items.append(OperationDetail( task_name=r[1] or "", product_sn=r[2] or "", material_name=r[3] or "", remark=r[4] or "", time=_to_iso(r[0]), )) elif action_type == "record": stmt = ( select(TaskRecord.created_at, Task.task_name, Product.serial_number, Product.material_name, TaskRecord.remark) .join(Task, TaskRecord.task_id == Task.id) .join(Product, Task.product_id == Product.id) .where( Task.assignee_id == user_id, or_(TaskRecord.remark.is_(None), ~TaskRecord.remark.like("[%")), ) .order_by(TaskRecord.created_at.desc()) ) if since: stmt = stmt.where(TaskRecord.created_at >= since) if until: stmt = stmt.where(TaskRecord.created_at <= until) rows = (await db.execute(stmt)).all() for r in rows: items.append(OperationDetail( task_name=r[1] or "", product_sn=r[2] or "", material_name=r[3] or "", remark=r[4] or "", time=_to_iso(r[0]), )) return items # ============================================================ # 人员负载(按人聚合在制品设备 — 独立「人员看板」) # ============================================================ async def get_people_workload(db: AsyncSession) -> list[PersonWorkload]: """上帝视角 — 按负责人聚合当前在制品设备(WIP/PENDING 任务,product 去重,含滞留时长)。""" from app.models.task import Task, TASK_STATUS_PENDING, TASK_STATUS_WIP from app.models.product import Product from app.models.holiday import Holiday from app.core.time_utils import get_beijing_time, to_beijing, working_duration_hours # 读取节假日(排除非工作日) hres = await db.execute(select(Holiday.day)) holidays = {r[0] for r in hres} stmt = ( select( Task.assignee_id, Product.id, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model, Task.status, Task.received_at, Task.created_at, ) .join(Product, Task.product_id == Product.id) .where( Task.status.in_([TASK_STATUS_PENDING, TASK_STATUS_WIP]), Task.assignee_id.isnot(None), ) .order_by(Task.assignee_id, Product.created_at.desc()) ) result = await db.execute(stmt) rows = result.all() def _to_bj(dt): return to_beijing(dt) now = get_beijing_time() # 中间聚合:by_assignee[assignee][product_id] -> 设备快照 + 最早接手时间 agg: dict[str, dict[str, dict]] = {} for row in rows: assignee = row[0] product_id = str(row[1]) status = row[6] or "" received_dt = _to_bj(row[7] or row[8]) entry = agg.setdefault(assignee, {}).setdefault(product_id, { "status": "", "earliest_dt": None, "serial": row[2] or "", "ext": row[3] or None, "mat": row[4] or "", "spec": row[5] or "", }) # 状态优先级 WIP > PENDING if status == "WIP": entry["status"] = "WIP" elif not entry["status"]: entry["status"] = status # 取最早接手时间 if received_dt and (entry["earliest_dt"] is None or received_dt < entry["earliest_dt"]): entry["earliest_dt"] = received_dt raw_ids = list(agg.keys()) name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(raw_ids) workloads: list[PersonWorkload] = [] for assignee, prods in agg.items(): devices: list[PersonDevice] = [] for product_id, e in prods.items(): dt = e["earliest_dt"] hours = working_duration_hours(dt, now, holidays) if dt else 0.0 received_str = dt.strftime("%m-%d %H:%M") if dt else None devices.append(PersonDevice( product_id=product_id, serial_number=e["serial"], external_serial=e["ext"], material_name=e["mat"], spec_model=e["spec"], task_status=e["status"], duration_hours=hours, received_at=received_str, )) workloads.append(PersonWorkload( assignee_id=assignee, assignee_name=name_map.get(assignee, assignee), device_count=len(devices), devices=devices, )) workloads.sort(key=lambda w: w.device_count, reverse=True) return workloads # ============================================================ # 历史人员看板(按人聚合历史备注记录,按时段过滤) # ============================================================ async def get_people_history( db: AsyncSession, since: datetime | None = None, until: datetime | None = None, assignee_id: str | None = None, spec_model: str | None = None, product_sn: str | None = None, task_name: str | None = None, ) -> list[PersonHistoryRecord]: """上帝视角 — 人员效能与工时台账(平铺 Task 明细,含 WIP/PENDING/COMPLETED)。""" from app.models.task import Task, TASK_STATUS_WIP, TASK_STATUS_PENDING, TASK_STATUS_COMPLETED from app.models.product import Product from app.core.time_utils import get_beijing_time, BEIJING_TZ, to_beijing, working_duration_hours from app.models.holiday import Holiday from sqlalchemy import func now = get_beijing_time() # 读取节假日(排除非工作日) hres = await db.execute(select(Holiday.day)) holidays = {r[0] for r in hres} # ── 关联 TaskRecord:最新有效备注 + 记录总数 ── from app.models.task import TaskRecord from sqlalchemy import and_, or_ # 系统自动备注关键字(转交/移交/撤回/派发/驳回等) sys_remark = or_( TaskRecord.remark.ilike("%转交%"), TaskRecord.remark.ilike("%移交%"), TaskRecord.remark.ilike("%撤回%"), TaskRecord.remark.ilike("%重新接手%"), TaskRecord.remark.ilike("%派发%"), TaskRecord.remark.ilike("%分配%"), TaskRecord.remark.ilike("%驳回%"), TaskRecord.remark.ilike("%返工%"), TaskRecord.remark.ilike("%完工%"), ) # 最新有效备注(排除系统备注,按时间倒序取第一条) valid_remark_subq = ( select( TaskRecord.task_id, TaskRecord.remark.label("latest_valid_remark"), func.row_number().over( partition_by=TaskRecord.task_id, order_by=TaskRecord.created_at.desc(), ).label("rn"), ) .where( TaskRecord.remark.isnot(None), func.trim(TaskRecord.remark) != "", ~sys_remark, ) ).subquery("vr") # 记录总数(含系统备注) record_count_subq = ( select( TaskRecord.task_id, func.count(TaskRecord.id).label("record_count"), ) .group_by(TaskRecord.task_id) ).subquery("rc") stmt = ( select( Task.id, Task.task_name, Task.assignee_id, Task.status, Task.received_at, Task.created_at, Task.completed_at, Product.serial_number, Product.external_serial, Product.material_name, Product.spec_model, valid_remark_subq.c.latest_valid_remark, func.coalesce(record_count_subq.c.record_count, 0).label("record_count"), ) .join(Product, Task.product_id == Product.id) .outerjoin( valid_remark_subq, and_(valid_remark_subq.c.task_id == Task.id, valid_remark_subq.c.rn == 1), ) .outerjoin(record_count_subq, record_count_subq.c.task_id == Task.id) .where( Task.status.in_([TASK_STATUS_WIP, TASK_STATUS_PENDING, TASK_STATUS_COMPLETED]), Task.assignee_id.isnot(None), ) ) # 多维筛选(模糊/精确匹配) if assignee_id: stmt = stmt.where(Task.assignee_id == assignee_id) if spec_model: stmt = stmt.where(Product.spec_model.ilike(f"%{spec_model.strip()}%")) if product_sn: stmt = stmt.where(Product.serial_number.ilike(f"%{product_sn.strip()}%")) if task_name: stmt = stmt.where(Task.task_name.ilike(f"%{task_name.strip()}%")) # 时间交集:start = COALESCE(received_at, created_at),end = COALESCE(completed_at, now) start_expr = func.coalesce(Task.received_at, Task.created_at) end_expr = func.coalesce(Task.completed_at, now) if until: stmt = stmt.where(start_expr <= until) if since: stmt = stmt.where(end_expr >= since) # 排序:先进行中(WIP/PENDING),后已完成(COMPLETED),状态内按开始时间降序 from sqlalchemy import case status_prio = case( (Task.status == TASK_STATUS_WIP, 0), (Task.status == TASK_STATUS_PENDING, 1), (Task.status == TASK_STATUS_COMPLETED, 2), else_=3, ) stmt = stmt.order_by( status_prio.asc(), func.coalesce(Task.received_at, Task.created_at).desc(), ) result = await db.execute(stmt) rows = result.all() def _to_bj(dt): if not dt: return None if dt.tzinfo is None: from datetime import timezone as dt_timezone dt = dt.replace(tzinfo=dt_timezone.utc).astimezone(BEIJING_TZ) else: dt = dt.astimezone(BEIJING_TZ) return dt raw_ids = list({r[2] for r in rows if r[2]}) name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(raw_ids) records: list[PersonHistoryRecord] = [] for row in rows: received_dt = _to_bj(row[4] or row[5]) # received_at or created_at completed_dt = _to_bj(row[6]) # completed_at(进行中为 None) if completed_dt: end_dt = completed_dt completed_str = completed_dt.strftime("%m-%d %H:%M") else: end_dt = now completed_str = None hours = working_duration_hours(received_dt, end_dt, holidays) if received_dt else 0.0 records.append(PersonHistoryRecord( task_id=str(row[0]), task_name=row[1] or "", assignee_id=row[2] or "", assignee_name=name_map.get(row[2] or "", row[2] or ""), product_sn=row[7] or "", external_serial=row[8] or None, material_name=row[9] or "", spec_model=row[10] or "", status=row[3] or "", received_at=received_dt.strftime("%m-%d %H:%M") if received_dt else None, completed_at=completed_str, duration_hours=hours, latest_valid_remark=row[11], record_count=row[12] or 0, )) return records # ============================================================ # 协同留言搜索(上帝视角 — 全厂) # ============================================================ async def search_product_messages( db: AsyncSession, keyword: str = "", skip: int = 0, limit: int = 50, ) -> ProductMessageList: """ 上帝视角 — 全厂所有产品的协同留言。 关联链: ProductMessage → Product → (material_name, serial_number) 搜索支持: SN码、物料名称、留言人 排序: created_at 倒序(最新在前) """ from app.models.message import ProductMessage from app.models.product import Product from app.core.time_utils import BEIJING_TZ # 基础查询 — 同时取出 16位追溯码 + 业务序列号 stmt = ( select(ProductMessage, Product.serial_number, Product.material_name, Product.external_serial) .join(Product, ProductMessage.product_id == Product.id) ) # 关键词搜索 if keyword and keyword.strip(): kw = f"%{keyword.strip()}%" stmt = stmt.where(or_( Product.serial_number.ilike(kw), Product.external_serial.ilike(kw), Product.material_name.ilike(kw), ProductMessage.operator_id.ilike(kw), ProductMessage.content.ilike(kw), )) # 总数 count_stmt = select(func.count()).select_from(stmt.subquery()) total = await db.scalar(count_stmt) or 0 # 分页 + 排序 stmt = stmt.order_by(ProductMessage.created_at.desc()).offset(skip).limit(limit) result = await db.execute(stmt) rows = result.all() # 收集 operator_id → 批量翻译中文姓名 raw_ids = list({row[0].operator_id for row in rows if row[0].operator_id}) name_map: dict[str, str] = {} if raw_ids: from app.services.mom_cache import get_display_names name_map = get_display_names(raw_ids) items: list[ProductMessageItem] = [] for msg, sn, mat_name, ext_sn in rows: t = msg.created_at if t: if t.tzinfo is None: # 旧数据(datetime.utcnow):naive UTC → 转北京时间 from datetime import timezone as dt_timezone t = t.replace(tzinfo=dt_timezone.utc).astimezone(BEIJING_TZ) else: t = t.astimezone(BEIJING_TZ) time_str = t.isoformat() if t else "" items.append(ProductMessageItem( id=str(msg.id), content=msg.content, operator_name=name_map.get(msg.operator_id, msg.operator_id), product_sn=sn or "", external_serial=ext_sn or None, material_name=mat_name or "", created_at=time_str, )) return ProductMessageList(items=items, total=total)