2026-08-14 10:19:44 +08:00
|
|
|
|
"""效能分析服务 — 个人能力图谱 + 设备流转对比(ECharts 数据源)
|
|
|
|
|
|
|
|
|
|
|
|
数据来源:
|
|
|
|
|
|
- Task 工序/任务节点,携带 assignee_id、received_at、completed_at → 计算单台耗时。
|
|
|
|
|
|
- Product 设备身份证(sn)、规格型号(spec_model)、物料名,用于筛选与展示。
|
|
|
|
|
|
- TaskRecord 备注时间线(本模块当前仅作耗时主数据补充,可按需在后续下钻中引入)。
|
|
|
|
|
|
|
|
|
|
|
|
耗时口径:
|
|
|
|
|
|
- 单台真实耗时 = (completed_at or now) - (received_at or created_at),单位小时。
|
|
|
|
|
|
- naive datetime 按 UTC 处理,统一换算北京时间。
|
|
|
|
|
|
"""
|
|
|
|
|
|
from datetime import datetime, timezone
|
|
|
|
|
|
|
|
|
|
|
|
from sqlalchemy import select, func
|
|
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
|
|
from pydantic import BaseModel
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# Schemas(与前端 src/services/analyticsApi.ts 对齐)
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
class CapabilityDataPoint(BaseModel):
|
|
|
|
|
|
value: float | None # 该人员在该设备上的总耗时(未触及为 None,真实 0 为 0.0)
|
|
|
|
|
|
spec_model: str # 规格型号
|
|
|
|
|
|
status: str # 该设备当前状态 WIP/PENDING/COMPLETED/—
|
|
|
|
|
|
first_received_at: str # 最早接收时间 YYYY-MM-DD HH:mm
|
|
|
|
|
|
last_completed_at: str # 最后完成时间 YYYY-MM-DD HH:mm(无则空串)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class CapabilityDevice(BaseModel):
|
|
|
|
|
|
product_sn: str # 身份证
|
|
|
|
|
|
external_serial: str | None # 产品序列号(业务序列号)
|
|
|
|
|
|
spec_model: str # 规格型号
|
2026-08-14 15:43:22 +08:00
|
|
|
|
material_name: str # 物料名称
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class CapabilitySeries(BaseModel):
|
|
|
|
|
|
name: str # 人员姓名
|
|
|
|
|
|
assignee_id: str # 人员ID
|
|
|
|
|
|
data: list[CapabilityDataPoint] # 与 categories 严格对齐,未触及设备补 0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class CapabilityResponse(BaseModel):
|
|
|
|
|
|
categories: list[str] # X 轴:设备身份证(按时间升序)
|
|
|
|
|
|
devices: list[CapabilityDevice] # 与 categories 对齐的设备元数据
|
|
|
|
|
|
series: list[CapabilitySeries]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class FlowDevice(BaseModel):
|
|
|
|
|
|
product_sn: str
|
|
|
|
|
|
external_serial: str | None
|
|
|
|
|
|
material_name: str
|
|
|
|
|
|
spec_model: str
|
2026-08-14 11:44:41 +08:00
|
|
|
|
lead_time: float # 设备生命周期总时长(小时)= max_end - min_start
|
|
|
|
|
|
started_at: str # 设备最早介入时间(T0),格式 MM-DD HH:mm
|
2026-08-28 15:28:56 +08:00
|
|
|
|
total_days: int = 0 # 设备生产总天数(自然天,自最早介入至今)
|
|
|
|
|
|
total_workdays: int = 0 # 设备生产总天数(工作日,排除周末/节假日)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class FlowSeries(BaseModel):
|
2026-08-14 11:44:41 +08:00
|
|
|
|
name: str # 人员姓名
|
|
|
|
|
|
data: list[list] # 每项 = [device_index, start_offset, end_offset, task_name, duration](小时)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class FlowResponse(BaseModel):
|
|
|
|
|
|
devices: list[FlowDevice]
|
|
|
|
|
|
series: list[FlowSeries]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class AssigneeOption(BaseModel):
|
|
|
|
|
|
id: str
|
|
|
|
|
|
name: str
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-14 14:30:35 +08:00
|
|
|
|
class SpecModelOption(BaseModel):
|
|
|
|
|
|
spec_model: str
|
|
|
|
|
|
material_name: str
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class DeviceOption(BaseModel):
|
|
|
|
|
|
product_sn: str
|
|
|
|
|
|
external_serial: str | None
|
|
|
|
|
|
material_name: str
|
|
|
|
|
|
spec_model: str
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-14 10:19:44 +08:00
|
|
|
|
class AnalyticsOptions(BaseModel):
|
|
|
|
|
|
assignees: list[AssigneeOption]
|
2026-08-14 14:30:35 +08:00
|
|
|
|
spec_models: list[SpecModelOption]
|
|
|
|
|
|
devices: list[DeviceOption]
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class DeviceRecord(BaseModel):
|
|
|
|
|
|
task_name: str # 工序名
|
|
|
|
|
|
assignee_name: str # 负责人姓名
|
|
|
|
|
|
status: str # 任务状态
|
|
|
|
|
|
remark: str | None # 备注
|
|
|
|
|
|
images: list[str] # 图片 URL 列表
|
|
|
|
|
|
created_at: str # 记录时间 ISO
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# 时间工具
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
def _to_bj(dt: datetime | None) -> datetime | None:
|
|
|
|
|
|
"""naive datetime 按 UTC 处理,统一换算为北京时间。"""
|
|
|
|
|
|
from app.core.time_utils import BEIJING_TZ
|
|
|
|
|
|
if not dt:
|
|
|
|
|
|
return None
|
|
|
|
|
|
if dt.tzinfo is None:
|
|
|
|
|
|
dt = dt.replace(tzinfo=timezone.utc).astimezone(BEIJING_TZ)
|
|
|
|
|
|
else:
|
|
|
|
|
|
dt = dt.astimezone(BEIJING_TZ)
|
|
|
|
|
|
return dt
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-28 15:08:10 +08:00
|
|
|
|
def _duration_hours(start: datetime | None, end: datetime, holidays: set = None) -> float:
|
|
|
|
|
|
"""计算单台耗时(工作小时,排除周末/节假日),无开始时间返回 0。"""
|
2026-08-14 10:19:44 +08:00
|
|
|
|
if not start:
|
|
|
|
|
|
return 0.0
|
2026-08-28 15:08:10 +08:00
|
|
|
|
from app.core.time_utils import working_duration_hours
|
|
|
|
|
|
return working_duration_hours(start, end, holidays)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# 个人能力图谱 — 单台设备耗时对比(分组柱状图,X=设备身份证)
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
async def get_capability_profile(
|
|
|
|
|
|
db: AsyncSession,
|
|
|
|
|
|
assignee_ids: list[str] | None = None,
|
|
|
|
|
|
spec_models: list[str] | None = None,
|
|
|
|
|
|
since: datetime | None = None,
|
|
|
|
|
|
until: datetime | None = None,
|
|
|
|
|
|
) -> CapabilityResponse:
|
|
|
|
|
|
"""
|
|
|
|
|
|
个人能力图谱(单机颗粒度):X 轴 = 设备身份证,每个负责人一条柱状 series。
|
|
|
|
|
|
按 (负责人, 设备) 分组,累加该人在该设备所有工序的耗时,
|
|
|
|
|
|
series.data 与 categories 严格对齐,未触及设备补 0。
|
|
|
|
|
|
耗时口径:(coalesce(completed_at, now) - coalesce(received_at, created_at))。
|
|
|
|
|
|
"""
|
|
|
|
|
|
from app.models.task import (
|
|
|
|
|
|
Task, TASK_STATUS_WIP, TASK_STATUS_PENDING, TASK_STATUS_COMPLETED,
|
|
|
|
|
|
)
|
|
|
|
|
|
from app.models.product import Product
|
2026-08-28 15:08:10 +08:00
|
|
|
|
from app.models.holiday import Holiday
|
2026-08-14 10:19:44 +08:00
|
|
|
|
from app.core.time_utils import get_beijing_time
|
|
|
|
|
|
|
|
|
|
|
|
now = get_beijing_time()
|
|
|
|
|
|
|
2026-08-28 15:08:10 +08:00
|
|
|
|
# 读取节假日(排除非工作日)
|
|
|
|
|
|
hres = await db.execute(select(Holiday.day))
|
|
|
|
|
|
holidays = {r[0] for r in hres}
|
|
|
|
|
|
|
2026-08-14 10:19:44 +08:00
|
|
|
|
stmt = (
|
|
|
|
|
|
select(
|
|
|
|
|
|
Task.assignee_id, Task.status,
|
|
|
|
|
|
Task.received_at, Task.created_at, Task.completed_at,
|
|
|
|
|
|
Product.serial_number, Product.external_serial, Product.spec_model,
|
2026-08-14 15:43:22 +08:00
|
|
|
|
Product.material_name,
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
|
|
|
|
|
.join(Product, Task.product_id == Product.id)
|
|
|
|
|
|
.where(
|
|
|
|
|
|
Task.status.in_([TASK_STATUS_WIP, TASK_STATUS_PENDING, TASK_STATUS_COMPLETED]),
|
|
|
|
|
|
Task.assignee_id.isnot(None),
|
|
|
|
|
|
)
|
|
|
|
|
|
)
|
|
|
|
|
|
if assignee_ids:
|
|
|
|
|
|
stmt = stmt.where(Task.assignee_id.in_(assignee_ids))
|
|
|
|
|
|
if spec_models:
|
|
|
|
|
|
stmt = stmt.where(Product.spec_model.in_(spec_models))
|
|
|
|
|
|
|
|
|
|
|
|
# 时间交集(与 people-history 口径一致)
|
|
|
|
|
|
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)
|
|
|
|
|
|
|
|
|
|
|
|
result = await db.execute(stmt)
|
|
|
|
|
|
rows = result.all()
|
|
|
|
|
|
|
|
|
|
|
|
# 设备元数据 + 时间基准;人员×设备 耗时聚合
|
|
|
|
|
|
device_meta: dict[str, dict] = {}
|
|
|
|
|
|
agg: dict[tuple[str, str], dict] = {}
|
|
|
|
|
|
for row in rows:
|
|
|
|
|
|
assignee = row[0]
|
|
|
|
|
|
status = row[1]
|
|
|
|
|
|
start = _to_bj(row[2] or row[3]) # received_at or created_at
|
|
|
|
|
|
end = _to_bj(row[4]) if row[4] else now # completed_at or now
|
|
|
|
|
|
sn = row[5]
|
|
|
|
|
|
ext = row[6]
|
|
|
|
|
|
spec = row[7] or "未知型号"
|
2026-08-14 15:43:22 +08:00
|
|
|
|
mat = row[8] or ""
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
2026-08-14 15:43:22 +08:00
|
|
|
|
meta = device_meta.setdefault(sn, {"external_serial": ext, "spec_model": spec, "material_name": mat, "earliest": None})
|
2026-08-14 10:19:44 +08:00
|
|
|
|
if start and (meta["earliest"] is None or start < meta["earliest"]):
|
|
|
|
|
|
meta["earliest"] = start
|
|
|
|
|
|
|
|
|
|
|
|
entry = agg.setdefault((assignee, sn), {
|
|
|
|
|
|
"hours": 0.0, "prio": 9,
|
|
|
|
|
|
"first_start": None, "last_completed": None,
|
|
|
|
|
|
})
|
2026-08-28 15:08:10 +08:00
|
|
|
|
entry["hours"] += _duration_hours(start, end, holidays)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
if start and (entry["first_start"] is None or start < entry["first_start"]):
|
|
|
|
|
|
entry["first_start"] = start
|
|
|
|
|
|
completed_dt = _to_bj(row[4]) if row[4] else None
|
|
|
|
|
|
if completed_dt and (entry["last_completed"] is None or completed_dt > entry["last_completed"]):
|
|
|
|
|
|
entry["last_completed"] = completed_dt
|
|
|
|
|
|
prio = 0 if status == TASK_STATUS_WIP else (1 if status == TASK_STATUS_PENDING else 2)
|
|
|
|
|
|
entry["prio"] = min(entry["prio"], prio)
|
|
|
|
|
|
|
|
|
|
|
|
if not agg:
|
|
|
|
|
|
return CapabilityResponse(categories=[], devices=[], series=[])
|
|
|
|
|
|
|
|
|
|
|
|
# X 轴:按最早接收/创建时间升序的设备身份证
|
|
|
|
|
|
categories = sorted(device_meta.keys(), key=lambda s: device_meta[s]["earliest"] or now)
|
|
|
|
|
|
devices = [
|
|
|
|
|
|
CapabilityDevice(
|
|
|
|
|
|
product_sn=sn,
|
|
|
|
|
|
external_serial=device_meta[sn]["external_serial"],
|
|
|
|
|
|
spec_model=device_meta[sn]["spec_model"],
|
2026-08-14 15:43:22 +08:00
|
|
|
|
material_name=device_meta[sn]["material_name"],
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
|
|
|
|
|
for sn in categories
|
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
|
|
# 人员姓名映射
|
|
|
|
|
|
raw_ids = list({a for (a, _) in 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)
|
|
|
|
|
|
|
|
|
|
|
|
STATUS_CODE = {0: "WIP", 1: "PENDING", 2: "COMPLETED"}
|
|
|
|
|
|
|
|
|
|
|
|
by_assignee: dict[str, dict[str, dict]] = {}
|
|
|
|
|
|
for (a, sn), entry in agg.items():
|
|
|
|
|
|
by_assignee.setdefault(a, {})[sn] = entry
|
|
|
|
|
|
|
|
|
|
|
|
series: list[CapabilitySeries] = []
|
|
|
|
|
|
for a in sorted(by_assignee.keys()):
|
|
|
|
|
|
sn_entries = by_assignee[a]
|
|
|
|
|
|
data: list[CapabilityDataPoint] = []
|
|
|
|
|
|
for sn in categories:
|
|
|
|
|
|
e = sn_entries.get(sn)
|
|
|
|
|
|
if e:
|
|
|
|
|
|
data.append(CapabilityDataPoint(
|
|
|
|
|
|
value=round(e["hours"], 1),
|
|
|
|
|
|
spec_model=device_meta[sn]["spec_model"],
|
|
|
|
|
|
status=STATUS_CODE[e["prio"]],
|
|
|
|
|
|
first_received_at=e["first_start"].strftime("%Y-%m-%d %H:%M") if e["first_start"] else "",
|
|
|
|
|
|
last_completed_at=e["last_completed"].strftime("%Y-%m-%d %H:%M") if e["last_completed"] else "",
|
|
|
|
|
|
))
|
|
|
|
|
|
else:
|
|
|
|
|
|
data.append(CapabilityDataPoint(
|
|
|
|
|
|
value=None,
|
|
|
|
|
|
spec_model=device_meta[sn]["spec_model"],
|
|
|
|
|
|
status="—",
|
|
|
|
|
|
first_received_at="",
|
|
|
|
|
|
last_completed_at="",
|
|
|
|
|
|
))
|
|
|
|
|
|
series.append(CapabilitySeries(
|
|
|
|
|
|
name=name_map.get(a, a),
|
|
|
|
|
|
assignee_id=a,
|
|
|
|
|
|
data=data,
|
|
|
|
|
|
))
|
|
|
|
|
|
series.sort(key=lambda s: s.name)
|
|
|
|
|
|
return CapabilityResponse(categories=categories, devices=devices, series=series)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# 流转对比 — 设备各工序耗时(堆叠柱状)
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
async def get_flow_compare(
|
|
|
|
|
|
db: AsyncSession,
|
|
|
|
|
|
product_sns: list[str] | None = None,
|
|
|
|
|
|
spec_models: list[str] | None = None,
|
2026-08-28 15:34:17 +08:00
|
|
|
|
mode: str = "natural",
|
2026-08-14 10:19:44 +08:00
|
|
|
|
) -> FlowResponse:
|
|
|
|
|
|
"""
|
2026-08-14 11:44:41 +08:00
|
|
|
|
查询每台设备上各操作人的任务时间区间,拼装为「生命周期时间轴」区间图:
|
|
|
|
|
|
x 轴 = 设备身份证,y 轴 = 相对该设备 T0(最早介入时间)的小时偏移。
|
|
|
|
|
|
每个任务一根悬空区间柱(start_offset -> end_offset),并行任务可并排显示。
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
设备来源:
|
|
|
|
|
|
- 传 product_sns:按给定身份证;
|
|
|
|
|
|
- 仅传 spec_models:这些型号下最近有流转的 20 台设备;
|
|
|
|
|
|
- 都未传:返回空。
|
|
|
|
|
|
|
2026-08-14 11:44:41 +08:00
|
|
|
|
data 每项 = [device_index, start_offset, end_offset, task_name, duration](单位小时)。
|
2026-08-14 10:19:44 +08:00
|
|
|
|
"""
|
|
|
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
|
now = get_beijing_time()
|
|
|
|
|
|
|
|
|
|
|
|
if product_sns:
|
|
|
|
|
|
product_rows = (await db.execute(
|
|
|
|
|
|
select(
|
|
|
|
|
|
Product.id, Product.serial_number, Product.external_serial,
|
|
|
|
|
|
Product.material_name, Product.spec_model,
|
|
|
|
|
|
).where(Product.serial_number.in_(product_sns))
|
|
|
|
|
|
)).all()
|
|
|
|
|
|
elif spec_models:
|
|
|
|
|
|
latest_subq = (
|
|
|
|
|
|
select(
|
|
|
|
|
|
Task.product_id,
|
|
|
|
|
|
func.max(func.coalesce(Task.received_at, Task.created_at)).label("latest"),
|
|
|
|
|
|
)
|
|
|
|
|
|
.group_by(Task.product_id)
|
|
|
|
|
|
.subquery()
|
|
|
|
|
|
)
|
|
|
|
|
|
product_rows = (await db.execute(
|
|
|
|
|
|
select(
|
|
|
|
|
|
Product.id, Product.serial_number, Product.external_serial,
|
|
|
|
|
|
Product.material_name, Product.spec_model,
|
|
|
|
|
|
)
|
|
|
|
|
|
.join(latest_subq, latest_subq.c.product_id == Product.id)
|
|
|
|
|
|
.where(Product.spec_model.in_(spec_models))
|
|
|
|
|
|
.order_by(latest_subq.c.latest.desc())
|
|
|
|
|
|
.limit(20)
|
|
|
|
|
|
)).all()
|
|
|
|
|
|
else:
|
|
|
|
|
|
return FlowResponse(devices=[], series=[])
|
|
|
|
|
|
|
|
|
|
|
|
if not product_rows:
|
|
|
|
|
|
return FlowResponse(devices=[], series=[])
|
|
|
|
|
|
|
|
|
|
|
|
devices = [
|
|
|
|
|
|
FlowDevice(
|
|
|
|
|
|
product_sn=r[1], external_serial=r[2],
|
|
|
|
|
|
material_name=r[3] or "", spec_model=r[4] or "",
|
2026-08-14 11:44:41 +08:00
|
|
|
|
lead_time=0.0,
|
|
|
|
|
|
started_at="",
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
|
|
|
|
|
for r in product_rows
|
|
|
|
|
|
]
|
|
|
|
|
|
# 显式传入身份证时按输入顺序排列
|
|
|
|
|
|
if product_sns:
|
|
|
|
|
|
order = {sn: i for i, sn in enumerate(product_sns)}
|
|
|
|
|
|
devices.sort(key=lambda d: order.get(d.product_sn, len(order)))
|
|
|
|
|
|
|
|
|
|
|
|
id_to_sn = {r[0]: r[1] for r in product_rows}
|
|
|
|
|
|
sn_to_index = {d.product_sn: i for i, d in enumerate(devices)}
|
|
|
|
|
|
product_ids = [r[0] for r in product_rows]
|
|
|
|
|
|
|
2026-08-14 11:44:41 +08:00
|
|
|
|
# ── 这些设备的全部任务(含工序名,用于区间图) ──
|
2026-08-14 10:19:44 +08:00
|
|
|
|
task_rows = (await db.execute(
|
|
|
|
|
|
select(
|
2026-08-14 11:44:41 +08:00
|
|
|
|
Task.product_id, Task.assignee_id, Task.task_name,
|
2026-08-14 10:19:44 +08:00
|
|
|
|
Task.received_at, Task.created_at, Task.completed_at,
|
2026-08-14 11:44:41 +08:00
|
|
|
|
Task.task_type,
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
|
|
|
|
|
.where(
|
|
|
|
|
|
Task.product_id.in_(product_ids),
|
|
|
|
|
|
Task.status.in_([TASK_STATUS_WIP, TASK_STATUS_PENDING, TASK_STATUS_COMPLETED]),
|
|
|
|
|
|
Task.assignee_id.isnot(None),
|
|
|
|
|
|
)
|
|
|
|
|
|
)).all()
|
|
|
|
|
|
|
2026-08-14 11:44:41 +08:00
|
|
|
|
# ── 解析为区间记录,并计算每台设备的 T0(最早)与 max_end(最晚) ──
|
|
|
|
|
|
intervals: list[dict] = []
|
|
|
|
|
|
t0_by_device: dict[int, datetime] = {}
|
|
|
|
|
|
max_end_by_device: dict[int, datetime] = {}
|
2026-08-14 10:19:44 +08:00
|
|
|
|
for row in task_rows:
|
2026-08-14 11:44:41 +08:00
|
|
|
|
pid, assignee_id, task_name, received_at, created_at, completed_at, task_type = row
|
2026-08-14 10:19:44 +08:00
|
|
|
|
sn = id_to_sn.get(pid)
|
|
|
|
|
|
if sn is None or sn not in sn_to_index:
|
|
|
|
|
|
continue
|
|
|
|
|
|
idx = sn_to_index[sn]
|
|
|
|
|
|
start = _to_bj(received_at or created_at)
|
|
|
|
|
|
end = _to_bj(completed_at) if completed_at else now
|
2026-08-14 11:44:41 +08:00
|
|
|
|
if start is None:
|
|
|
|
|
|
continue
|
|
|
|
|
|
is_main = 0 if task_type == "SPAWN" else 1
|
|
|
|
|
|
intervals.append({
|
|
|
|
|
|
"idx": idx,
|
|
|
|
|
|
"assignee_id": assignee_id,
|
|
|
|
|
|
"task_name": (task_name or "").strip() or "未命名工序",
|
|
|
|
|
|
"start": start,
|
|
|
|
|
|
"end": end,
|
|
|
|
|
|
"is_main": is_main,
|
|
|
|
|
|
})
|
|
|
|
|
|
if idx not in t0_by_device or start < t0_by_device[idx]:
|
|
|
|
|
|
t0_by_device[idx] = start
|
|
|
|
|
|
if idx not in max_end_by_device or end > max_end_by_device[idx]:
|
|
|
|
|
|
max_end_by_device[idx] = end
|
|
|
|
|
|
|
|
|
|
|
|
# ── 计算每台设备 lead_time(最大 end - 最小 start,小时),重建 devices ──
|
|
|
|
|
|
lead_time_by_index = {
|
|
|
|
|
|
idx: round((max_end_by_device[idx] - t0_by_device[idx]).total_seconds() / 3600, 1)
|
|
|
|
|
|
for idx in t0_by_device
|
|
|
|
|
|
}
|
2026-08-28 15:28:56 +08:00
|
|
|
|
# 🔧 设备生产总天数(自然天 + 工作日,自最早介入至今)
|
|
|
|
|
|
import math
|
|
|
|
|
|
from app.core.time_utils import to_beijing as _tb, working_duration_hours as _wdh
|
|
|
|
|
|
from app.models.holiday import Holiday as _Holiday
|
|
|
|
|
|
hres = await db.execute(select(_Holiday.day))
|
|
|
|
|
|
_holidays = {r[0] for r in hres}
|
|
|
|
|
|
|
|
|
|
|
|
def _total_days(t0):
|
|
|
|
|
|
t0_bj = _tb(t0)
|
|
|
|
|
|
if not t0_bj:
|
|
|
|
|
|
return 1, 1
|
|
|
|
|
|
natural = max(1, math.ceil((now - t0_bj).total_seconds() / 86400))
|
|
|
|
|
|
workdays = max(1, math.ceil(_wdh(t0_bj, now, _holidays) / 24))
|
|
|
|
|
|
return natural, workdays
|
|
|
|
|
|
|
|
|
|
|
|
days_map = {i: _total_days(t0_by_device[i]) for i in t0_by_device}
|
2026-08-14 11:44:41 +08:00
|
|
|
|
devices = [
|
|
|
|
|
|
FlowDevice(
|
|
|
|
|
|
product_sn=d.product_sn, external_serial=d.external_serial,
|
|
|
|
|
|
material_name=d.material_name, spec_model=d.spec_model,
|
|
|
|
|
|
lead_time=lead_time_by_index.get(i, 0.0),
|
|
|
|
|
|
started_at=t0_by_device[i].strftime("%m-%d %H:%M") if i in t0_by_device else "",
|
2026-08-28 15:28:56 +08:00
|
|
|
|
total_days=days_map.get(i, (1, 1))[0],
|
|
|
|
|
|
total_workdays=days_map.get(i, (1, 1))[1],
|
2026-08-14 11:44:41 +08:00
|
|
|
|
)
|
|
|
|
|
|
for i, d in enumerate(devices)
|
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
|
|
# ── 按人分组,转为相对 T0 的小时偏移区间 ──
|
|
|
|
|
|
by_assignee: dict[str, list[list]] = {}
|
2026-08-28 15:34:17 +08:00
|
|
|
|
# 工作日模式:把时间偏移换算为「排除周末/节假日的工作小时」,压缩休息日
|
|
|
|
|
|
if mode == "workdays":
|
|
|
|
|
|
from app.core.time_utils import working_duration_hours as _wdh
|
|
|
|
|
|
from app.models.holiday import Holiday as _Holiday
|
|
|
|
|
|
hres = await db.execute(select(_Holiday.day))
|
|
|
|
|
|
_holidays = {r[0] for r in hres}
|
2026-08-14 11:44:41 +08:00
|
|
|
|
for it in intervals:
|
|
|
|
|
|
t0 = t0_by_device[it["idx"]]
|
2026-08-28 15:34:17 +08:00
|
|
|
|
if mode == "workdays":
|
|
|
|
|
|
start_offset = _wdh(t0, it["start"], _holidays)
|
|
|
|
|
|
end_offset = _wdh(t0, it["end"], _holidays)
|
|
|
|
|
|
else:
|
|
|
|
|
|
start_offset = round((it["start"] - t0).total_seconds() / 3600, 1)
|
|
|
|
|
|
end_offset = round((it["end"] - t0).total_seconds() / 3600, 1)
|
2026-08-14 11:44:41 +08:00
|
|
|
|
duration = round(end_offset - start_offset, 1)
|
|
|
|
|
|
by_assignee.setdefault(it["assignee_id"], []).append(
|
|
|
|
|
|
[it["idx"], start_offset, end_offset, it["task_name"], duration, it["is_main"]]
|
|
|
|
|
|
)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
|
|
|
|
|
# 翻译人员姓名
|
2026-08-14 11:44:41 +08:00
|
|
|
|
raw_ids = list(by_assignee.keys())
|
2026-08-14 10:19:44 +08:00
|
|
|
|
name_map: dict[str, str] = {}
|
|
|
|
|
|
if raw_ids:
|
|
|
|
|
|
from app.services.mom_cache import get_display_names
|
|
|
|
|
|
name_map = get_display_names(raw_ids)
|
|
|
|
|
|
|
|
|
|
|
|
series = [
|
2026-08-14 11:44:41 +08:00
|
|
|
|
FlowSeries(name=name_map.get(a, a), data=by_assignee[a])
|
|
|
|
|
|
for a in sorted(by_assignee.keys())
|
2026-08-14 10:19:44 +08:00
|
|
|
|
]
|
|
|
|
|
|
return FlowResponse(devices=devices, series=series)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# 设备备注记录 — 单台设备的全部备注/照片(弹窗下钻)
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
async def get_device_records(
|
|
|
|
|
|
db: AsyncSession,
|
|
|
|
|
|
product_sn: str,
|
|
|
|
|
|
assignee_ids: list[str] | None = None,
|
|
|
|
|
|
) -> list[DeviceRecord]:
|
|
|
|
|
|
"""查询某台设备(身份证)的所有任务备注记录,含图片,按时间倒序。
|
|
|
|
|
|
可选按负责人过滤(assignee_ids)。"""
|
|
|
|
|
|
import json
|
|
|
|
|
|
from app.models.task import Task, TaskRecord
|
|
|
|
|
|
from app.models.product import Product
|
|
|
|
|
|
|
|
|
|
|
|
product_id = await db.scalar(
|
|
|
|
|
|
select(Product.id).where(Product.serial_number == product_sn)
|
|
|
|
|
|
)
|
|
|
|
|
|
if not product_id:
|
|
|
|
|
|
return []
|
|
|
|
|
|
|
|
|
|
|
|
stmt = (
|
|
|
|
|
|
select(
|
|
|
|
|
|
TaskRecord.remark, TaskRecord.images, TaskRecord.created_at,
|
|
|
|
|
|
Task.task_name, Task.assignee_id, Task.status,
|
|
|
|
|
|
)
|
|
|
|
|
|
.join(Task, TaskRecord.task_id == Task.id)
|
|
|
|
|
|
.where(Task.product_id == product_id)
|
|
|
|
|
|
)
|
|
|
|
|
|
if assignee_ids:
|
|
|
|
|
|
stmt = stmt.where(Task.assignee_id.in_(assignee_ids))
|
|
|
|
|
|
stmt = stmt.order_by(TaskRecord.created_at.desc())
|
|
|
|
|
|
result = await db.execute(stmt)
|
|
|
|
|
|
rows = result.all()
|
|
|
|
|
|
|
|
|
|
|
|
raw_ids = list({r[4] for r in rows if r[4]})
|
|
|
|
|
|
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[DeviceRecord] = []
|
|
|
|
|
|
for row in rows:
|
|
|
|
|
|
remark, images_raw, created, task_name, assignee, status = row
|
|
|
|
|
|
try:
|
|
|
|
|
|
images = json.loads(images_raw) if images_raw else []
|
|
|
|
|
|
except (json.JSONDecodeError, TypeError):
|
|
|
|
|
|
images = []
|
|
|
|
|
|
if not isinstance(images, list):
|
|
|
|
|
|
images = []
|
|
|
|
|
|
records.append(DeviceRecord(
|
|
|
|
|
|
task_name=task_name or "",
|
|
|
|
|
|
assignee_name=name_map.get(assignee or "", assignee or ""),
|
|
|
|
|
|
status=status or "",
|
|
|
|
|
|
remark=remark,
|
|
|
|
|
|
images=images,
|
|
|
|
|
|
created_at=_to_bj(created).isoformat() if created else "",
|
|
|
|
|
|
))
|
|
|
|
|
|
return records
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
# 下拉选项 — 负责人 + 规格型号
|
|
|
|
|
|
# ============================================================
|
|
|
|
|
|
|
|
|
|
|
|
async def get_analytics_options(
|
|
|
|
|
|
db: AsyncSession,
|
|
|
|
|
|
assignee_ids: list[str] | None = None,
|
|
|
|
|
|
spec_models: list[str] | None = None,
|
|
|
|
|
|
) -> AnalyticsOptions:
|
|
|
|
|
|
"""返回筛选栏选项,支持动态联动:
|
|
|
|
|
|
- 传入 spec_models:只返回碰过这些型号的人;
|
2026-08-14 14:30:35 +08:00
|
|
|
|
- 传入 assignee_ids:只返回这些人处理过的型号;
|
|
|
|
|
|
- devices:根据筛选条件返回关联设备(无筛选则返回最近流转的设备)。"""
|
2026-08-14 10:19:44 +08:00
|
|
|
|
from app.models.task import Task
|
|
|
|
|
|
from app.models.product import Product
|
|
|
|
|
|
|
2026-08-14 14:30:35 +08:00
|
|
|
|
# ── 负责人:从 tasks 去重(可选按 spec_models 过滤) ──
|
2026-08-14 10:19:44 +08:00
|
|
|
|
if spec_models:
|
|
|
|
|
|
assignee_stmt = (
|
|
|
|
|
|
select(Task.assignee_id)
|
|
|
|
|
|
.join(Product, Task.product_id == Product.id)
|
|
|
|
|
|
.where(
|
|
|
|
|
|
Task.assignee_id.isnot(None),
|
|
|
|
|
|
Product.spec_model.in_(spec_models),
|
|
|
|
|
|
)
|
|
|
|
|
|
.distinct()
|
|
|
|
|
|
)
|
|
|
|
|
|
else:
|
|
|
|
|
|
assignee_stmt = (
|
|
|
|
|
|
select(Task.assignee_id)
|
|
|
|
|
|
.where(Task.assignee_id.isnot(None))
|
|
|
|
|
|
.distinct()
|
|
|
|
|
|
)
|
|
|
|
|
|
assignee_rows = (await db.execute(assignee_stmt)).all()
|
|
|
|
|
|
distinct_assignees = sorted({r[0] for r in assignee_rows if r[0]})
|
|
|
|
|
|
|
|
|
|
|
|
name_map: dict[str, str] = {}
|
|
|
|
|
|
if distinct_assignees:
|
|
|
|
|
|
from app.services.mom_cache import get_display_names
|
|
|
|
|
|
name_map = get_display_names(distinct_assignees)
|
|
|
|
|
|
|
|
|
|
|
|
assignees = [
|
|
|
|
|
|
AssigneeOption(id=aid, name=name_map.get(aid, aid))
|
|
|
|
|
|
for aid in distinct_assignees
|
|
|
|
|
|
]
|
|
|
|
|
|
|
2026-08-14 14:30:35 +08:00
|
|
|
|
# ── 规格型号 + 物料名(可选按 assignee_ids 过滤;同型号取首个非空物料名) ──
|
2026-08-14 10:19:44 +08:00
|
|
|
|
if assignee_ids:
|
|
|
|
|
|
spec_stmt = (
|
2026-08-14 14:30:35 +08:00
|
|
|
|
select(Product.spec_model, Product.material_name)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
.join(Task, Task.product_id == Product.id)
|
|
|
|
|
|
.where(
|
|
|
|
|
|
Product.spec_model.isnot(None),
|
|
|
|
|
|
func.trim(Product.spec_model) != "",
|
|
|
|
|
|
Task.assignee_id.in_(assignee_ids),
|
|
|
|
|
|
)
|
2026-08-14 14:30:35 +08:00
|
|
|
|
.distinct()
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
|
|
|
|
|
else:
|
|
|
|
|
|
spec_stmt = (
|
2026-08-14 14:30:35 +08:00
|
|
|
|
select(Product.spec_model, Product.material_name)
|
2026-08-14 10:19:44 +08:00
|
|
|
|
.where(Product.spec_model.isnot(None), func.trim(Product.spec_model) != "")
|
2026-08-14 14:30:35 +08:00
|
|
|
|
.distinct()
|
|
|
|
|
|
)
|
|
|
|
|
|
spec_rows = (await db.execute(spec_stmt)).all()
|
|
|
|
|
|
spec_map: dict[str, str] = {}
|
|
|
|
|
|
for row in spec_rows:
|
|
|
|
|
|
sm = row[0]
|
|
|
|
|
|
mn = (row[1] or "").strip()
|
|
|
|
|
|
if sm not in spec_map:
|
|
|
|
|
|
spec_map[sm] = mn
|
|
|
|
|
|
elif mn and not spec_map[sm]:
|
|
|
|
|
|
spec_map[sm] = mn
|
|
|
|
|
|
spec_model_opts = [
|
|
|
|
|
|
SpecModelOption(spec_model=sm, material_name=spec_map[sm])
|
|
|
|
|
|
for sm in sorted(spec_map.keys())
|
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
|
|
# ── 设备字典(可选按 assignee_ids / spec_models 过滤;无筛选返回最近流转 100 台) ──
|
|
|
|
|
|
if assignee_ids or spec_models:
|
|
|
|
|
|
device_stmt = select(
|
|
|
|
|
|
Product.serial_number, Product.external_serial,
|
|
|
|
|
|
Product.material_name, Product.spec_model,
|
2026-08-14 10:19:44 +08:00
|
|
|
|
)
|
2026-08-14 14:30:35 +08:00
|
|
|
|
if assignee_ids:
|
|
|
|
|
|
device_stmt = device_stmt.join(
|
|
|
|
|
|
Task, Task.product_id == Product.id
|
|
|
|
|
|
).where(Task.assignee_id.in_(assignee_ids))
|
|
|
|
|
|
if spec_models:
|
|
|
|
|
|
device_stmt = device_stmt.where(Product.spec_model.in_(spec_models))
|
|
|
|
|
|
device_stmt = device_stmt.distinct()
|
|
|
|
|
|
else:
|
|
|
|
|
|
latest_subq = (
|
|
|
|
|
|
select(
|
|
|
|
|
|
Task.product_id,
|
|
|
|
|
|
func.max(func.coalesce(Task.received_at, Task.created_at)).label("latest"),
|
|
|
|
|
|
)
|
|
|
|
|
|
.group_by(Task.product_id)
|
|
|
|
|
|
.subquery()
|
|
|
|
|
|
)
|
|
|
|
|
|
device_stmt = (
|
|
|
|
|
|
select(
|
|
|
|
|
|
Product.serial_number, Product.external_serial,
|
|
|
|
|
|
Product.material_name, Product.spec_model,
|
|
|
|
|
|
)
|
|
|
|
|
|
.join(latest_subq, latest_subq.c.product_id == Product.id)
|
|
|
|
|
|
.order_by(latest_subq.c.latest.desc())
|
|
|
|
|
|
.limit(100)
|
|
|
|
|
|
)
|
|
|
|
|
|
device_rows = (await db.execute(device_stmt)).all()
|
|
|
|
|
|
devices = [
|
|
|
|
|
|
DeviceOption(
|
|
|
|
|
|
product_sn=r[0] or "",
|
|
|
|
|
|
external_serial=r[1] or None,
|
|
|
|
|
|
material_name=r[2] or "",
|
|
|
|
|
|
spec_model=r[3] or "",
|
|
|
|
|
|
)
|
|
|
|
|
|
for r in device_rows
|
|
|
|
|
|
]
|
2026-08-14 10:19:44 +08:00
|
|
|
|
|
2026-08-14 14:30:35 +08:00
|
|
|
|
return AnalyticsOptions(assignees=assignees, spec_models=spec_model_opts, devices=devices)
|