Files
track/backend/app/services/dashboard_service.py
duxingchen c3e03bb177 fix(透视表): 不再把任意已完成工序强制归在库
- 设备当前工序 = 最新主任务的工序名(task_name),不区分状态
- 「在库」只在该设备最新主任务真为在库工序时出现
- 完成的测试/生产等工序按原工序显示,不再误归在库
2026-08-28 16:57:02 +08:00

1132 lines
44 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
# ============================================================
# 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 # 驳回时间 ISOBEIJING_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 # 操作时间 ISOBEIJING_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=人员 或 工序,单元格=设备数量。
核心口径:**每台设备只统计一次**,按它「当前所处工序」归属——
1. 取该设备最新的一条主分支任务parent_task_id IS NULL 或 TRANSFER/RECOVERY
2. 设备当前工序 = 该最新主任务的工序名task_name不区分状态
- 「待确认」(PENDING) = 别人转给我但未接收
- 「在库」(COMPLETED) = 真正入库/生产完成
- 其他已完成工序(如「测试」完成)按原工序显示,不强制归「在库」
这样一台设备在「上一步已完成 + 下一步待确认」时只算一次(待确认),不会重复计数。
since/until 按设备最新主任务的创建时间过滤。
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,
Task.task_name,
Task.assignee_id,
Task.created_at,
)
.join(Task, Task.product_id == Product.id)
.where(
or_(
Task.parent_task_id.is_(None),
Task.task_type.in_(["TRANSFER", "RECOVERY"]),
)
)
.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, task_name, assignee, created in rows:
if pid in seen:
continue
seen.add(pid)
# 时间筛选:设备最新主任务的创建时间
if created is not None and created.tzinfo is None:
created = created.replace(tzinfo=dt_timezone.utc)
if since and created is not None and created < since:
continue
if until and created is not None and created > until:
continue
if dimension == "task_name":
# 设备当前工序 = 最新主任务的工序名(不区分状态,不强制完成态归在库)
key = task_name or ""
else:
key = assignee or "未分配"
device_cur[pid] = (spec or "未知型号", key, assignee or "")
# 聚合:规格 × 当前工序 → 设备数;同时收集负责人
agg: dict[tuple, int] = {}
assignee_map: dict[tuple, set] = {}
for spec, key, assignee in device_cur.values():
k = (spec, key)
agg[k] = agg.get(k, 0) + 1
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,
dimension_key=dim_display,
count=cnt,
assignees=assignees,
))
items.sort(key=lambda x: (x.spec_model, x.dimension_key))
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.utcnownaive 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)