Files
track/backend/app/services/dashboard_service.py
duxingchen f22eae315f fix: 售后回流口径统一 — 状态双字段同步、工序名归一、标签配色
产品存在 overall_status 与 status 两个状态字段,此前各写入点各写一份映射、
甚至只改 overall_status 不改 status,导致回流设备(product.status 停留在
OUTBOUND)污染看板统计口径。

- lifecycle.py: 把映射表收敛为单一来源 overall_to_product_status /
  sync_product_status;新增 normalize_after_sales_step,将售后设备沿用生产
  阶段写法的历史工序名(测试/维修)折算到售后区独立工序名
- product_service.py: 删除本地 _OVERALL_TO_STATUS 副本,改用共享函数
- dashboard_service.py: WIP 矩阵补出 lifecycle_phase 列,活跃任务判定
  (is_active) 提前到所有终结态判定之前,避免残留 OUTBOUND 被误判为完结
- scripts/fix_product_status.py: 历史数据修复脚本(一次性)
- constants/task.ts: 售后工序标签由红色改紫色 —— 红色在本系统是「驳回/危险」
  语义色,售后只是另一条流转支线,用红色会让操作员误以为设备报错
2026-09-15 10:57:29 +08:00

1428 lines
58 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Dashboard 统计服务 — 上帝视角(全厂全系统数据,不按用户过滤)"""
from datetime import datetime
from sqlalchemy import select, func, or_
from sqlalchemy.ext.asyncio import AsyncSession
from pydantic import BaseModel
from app.core.lifecycle import normalize_after_sales_step
# ============================================================
# 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
# 🔧 生命周期阶段:PRODUCTION(生产制造/发货测试) | AFTER_SALES(出库后返厂售后维修)
lifecycle_phase: str = "PRODUCTION"
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.status.label("task_status"),
Task.created_at,
Task.completed_at,
Product.current_location_id,
Product.overall_status,
Product.status,
Product.lifecycle_phase,
)
.outerjoin(Task, Task.product_id == Product.id)
.where(
or_(
Task.id.is_(None), # 🔧 允许该产品完全没有任务记录(新建档/待接收)
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, tstatus, created, completed, loc, overall_status, product_status, lifecycle 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
# ── 活跃任务判定:决定"按当前工序归列"还是"降级判完结态"的唯一依据 ──
# 必须在所有终结态判定之前。否则「已出库」的残留状态(设备回流后
# product.status 仍是 OUTBOUND,因为 receive_task/transfer_task 只同步
# overall_status 不同步 status)会把正在做售后工序的设备误判为终结态。
is_active = tstatus in ("WIP", "PENDING")
# ── 动静分离:仅终结(收口/在库/出库)设备按 since~until 过滤;在制设备无条件全量 ──
# 有活跃任务的一律视为在制,不做时间过滤(与"在制无条件全量"口径一致)
is_terminal = (not is_active) and (
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":
# ══ ① 第一优先级:有活跃任务(WIP/PENDING) → 一律按当前工序归列 ══
# 绝不看 overall_status / product.status 的终结态字段。这两个字段
# 在售后回流场景下会残留旧值(status 停在 OUTBOUND),若先判终结态,
# 正在做「发货测试」的设备会被吞进「已出库」列。
if is_active:
if tstatus == "PENDING":
# 转交后待接收(工序名还是占位符「待确认」)→ 统一归「待接收」
key = "待接收"
else:
key = task_name if (task_name and task_name != "—") else "待接收"
# 售后回流设备按折算后的工序名落列:发货测试 / 售后维修
key = normalize_after_sales_step(lifecycle, key)
# ══ ② 无活跃任务,才降级判完结态(多维防御,不依赖单一 loc / overall)══
# 已入库: 整体实收(已入库/在库) 或 status==ARCHIVED 或 收口节点「扫码入库」
elif 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:
# 🔧 对齐全景/产品管理:无任何任务的新品统一归「待接收」
if tstatus is None and task_name is None:
key = "待接收"
else:
# 已完工但未收口的工序(如「测试」完成),按原工序名显示
key = task_name if (task_name and task_name != "—") else "待接收"
# 🔧 售后回流设备的历史工序归一:老数据里售后设备的工序可能仍叫
# 「测试 / 维修」,折算进售后区的独立列,避免混入生产区同名工序
key = normalize_after_sales_step(lifecycle, key)
else:
key = assignee or "未分配"
device_cur[pid] = (spec or "未知型号", key, assignee or "", product_name or "")
from collections import Counter
# 聚合:规格 × 当前工序 → 设备数;负责人按人头计数,空值显式计为「未分配」
agg: dict[tuple, int] = {}
assignee_counter: dict[tuple, Counter] = {}
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)
raw = assignee if assignee and str(assignee).strip() not in ("", "—", "-", "null", "None") else ""
assignee_counter.setdefault(k, Counter())[raw] += 1
# 负责人 ID → 中文名
raw_ids: set[str] = set()
for counter in assignee_counter.values():
raw_ids |= set(counter.keys())
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))
def _disp(raw: str) -> str:
"""负责人显示名:空值统一为 未分配"""
return "未分配" if not raw else name_map.get(raw, raw)
def _fmt_label(raw: str, n: int) -> str:
"""负责人标签:>1 台必须带数量后缀;1 台不带"""
disp = _disp(raw)
return f"{disp}({n})" if n > 1 else disp
items: list[WipMatrixRow] = []
for (spec, key), cnt in agg.items():
dim_display = _disp(key) if dimension == "assignee" else key
counter = assignee_counter.get((spec, key), Counter())
# 负责人名单:未分配固定排最前,其余按台数降序、名称升序;带数量后缀(>1 台)
labels = [
_fmt_label(raw, n)
for raw, n in sorted(
counter.items(),
key=lambda kv: (0 if kv[0] == "" else 1, -kv[1], _disp(kv[0])),
)
]
items.append(WipMatrixRow(
spec_model=spec,
product_name=product_name_map.get((spec, key), ""),
dimension_key=dim_display,
count=cnt,
assignees=labels,
))
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,
Product.lifecycle_phase,
Task.task_name,
Task.assignee_id,
Task.status.label("task_status"),
Task.created_at,
Task.completed_at,
Task.received_at,
)
.outerjoin(Task, Task.product_id == Product.id)
.where(
or_(
Task.id.is_(None), # 🔧 允许该产品完全没有任务记录(新建档/待接收)
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, lifecycle, 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
# ── 活跃任务判定:与 get_wip_matrix 同口径,必须先于一切终结态判定 ──
# 否则「已出库」残留状态(回流后 pstatus 仍是 OUTBOUND)会把正在做
# 售后工序的设备误判为终结态,既错配时间过滤又错列。
is_active = tstatus in ("WIP", "PENDING")
# ── 动静分离:仅终结(收口/在库/出库)设备按 since~until 过滤;在制设备无条件全量 ──
is_terminal = (not is_active) and (
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 is_active:
if tstatus == "PENDING":
key = "待接收"
else:
key = task_name if (task_name and task_name != "—") else "待接收"
key = normalize_after_sales_step(lifecycle, key)
# ② 无活跃任务,才降级判完结态
elif 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:
# 🔧 对齐全景/产品管理:无任何任务的新品统一归「待接收」
if tstatus is None and task_name is None:
key = "待接收"
else:
key = task_name if (task_name and task_name != "—") else "待接收"
# 🔧 必须与 get_wip_matrix 同口径归一:否则矩阵格子里有数、
# 点进来下钻却是空的(售后设备的历史「测试/维修」工序)
key = normalize_after_sales_step(lifecycle, key)
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 "",
lifecycle_phase=lifecycle or "PRODUCTION",
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)