"""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 轴) dimension_key: str # 人员姓名 或 工序名称(X 轴) count: int = 0 # 该交叉点的设备数量 assignees: list[str] = [] # 该交叉点涉及的主负责人(中文名,去重) 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 # ── 产品(实时快照,不过滤) ── p_total = await db.scalar(select(func.count(Product.id))) p_pending = await db.scalar(select(func.count(Product.id)).where(Product.status == "pending")) p_progress = await db.scalar(select(func.count(Product.id)).where(Product.status == "in_progress")) p_done = await db.scalar(select(func.count(Product.id)).where(Product.status == "completed")) # ── 任务实时快照(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=人员 或 工序,单元格=设备数量。 覆盖全部任务状态(含已完成/在库),工序分布能看到「在库」(生产完成)。 只统计主分支(parent_task_id IS NULL 或 TRANSFER/RECOVERY),避免协助分支/返工重复计数。 since/until 按任务接手/创建时间过滤。 dimension: - assignee: 按任务负责人聚合(dimension_key 为中文姓名) - task_name: 按工序名聚合(dimension_key 为工序名,附主负责人) """ from app.models.task import Task from app.models.product import Product dim_expr = Task.task_name if dimension == "task_name" else Task.assignee_id stmt = ( select( Product.spec_model, dim_expr, func.count(Product.id), func.array_agg(func.distinct(Task.assignee_id)), ) .join(Task, Task.product_id == Product.id) # 🔧 只统计主分支(主线任务),排除 SPAWN 协助分支 .where( or_( Task.parent_task_id.is_(None), Task.task_type.in_(["TRANSFER", "RECOVERY"]), ) ) .group_by(Product.spec_model, dim_expr) .order_by(Product.spec_model, dim_expr) ) if since: stmt = stmt.where(func.coalesce(Task.received_at, Task.created_at) >= since) if until: stmt = stmt.where(func.coalesce(Task.received_at, Task.created_at) <= until) result = await db.execute(stmt) rows = result.all() # 收集负责人 ID → 中文名 raw_ids: set[str] = set() for r in rows: for aid in (r[3] or []): if aid: raw_ids.add(aid) 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, dim_key, cnt, assignee_ids in rows: dim_display = name_map.get(dim_key or "", dim_key or "未分配") if dimension == "assignee" else (dim_key or "—") assignees = [name_map.get(a, a) for a in (assignee_ids or []) if a] or [] items.append(WipMatrixRow( spec_model=spec or "未知型号", dimension_key=dim_display, count=cnt or 0, assignees=assignees, )) return items # ============================================================ # 人员操作明细(点击数字下钻 — 接收/转交/上传备注) # ============================================================ 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)