From ded5dfd7b325a33ed7e01f50bb52af61bb2bb588 Mon Sep 17 00:00:00 2001 From: yueli Date: Wed, 23 Sep 2026 17:16:26 +0800 Subject: [PATCH] =?UTF-8?q?feat(report):=20MOM=20=E7=B3=BB=E7=BB=9F?= =?UTF-8?q?=E6=97=A5=E6=8A=A5=EF=BC=88=E7=BA=AF=E6=A8=A1=E6=9D=BF=EF=BC=8C?= =?UTF-8?q?=E4=B8=8D=E6=8E=A5=20AI=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 每天 17:30(北京时间)自动汇总当日 新增/修改/删除、出库、借库, 渲染为纯文本邮件发送。 为什么不用 AI -------------- 此前那份日报由 Dify(外部 AI 服务,见 api/v1/ai_proxy.py)生成, 好处是能写人话,代价是: · 外部依赖一断,整份日报就没了 —— 这正是"AI 坏了"之后发生的事; · 按次计费、走网络、要管密钥; · **会在报表里编数字**,而日报是给人做决策用的,数字必须可信。 日报是固定格式的结构化统计(计数 + 清单),本来也不需要创造力。 数据全在库里:audit_logs / trans_outbound / trans_borrow。 口径与审计页保持一致 -------------------- · action 经 canon_action() 归一化(历史数据里 CREATE 与「新增」混用, 不归一会漏掉早期数据); · 默认排除 system 占位账号(与审计页默认的"真实用户"视图一致)。 · 凡有上限处一律显式写明「另有 N 条」,不静默截断。 · 变更值全部翻译成人类可读:状态码按模块查表、审批人/物料 ID → 名称、 布尔 → 是/否、人名串规整。翻不动就原样显示 —— 报表里的错译比裸数字 更难被发现,也更危险。 ★ 顺带修掉两个既有缺陷(日报每天都会走这条路径,不修会一直中招): 1) 定时任务被 8 个 worker 重复执行 gunicorn.conf.py 配了 8 worker,而 run.py 在**模块级**启动 APScheduler —— gunicorn 未开 preload_app,每个 worker 独立 import 一次 run.py, 于是每个 worker 各起一份调度器,同一 cron 任务被并发执行 8 次。 现有库存预警邮件因此每天重复发 8 封。 新增 app/utils/job_lock.py:用 PostgreSQL 会话级咨询锁做跨进程互斥, 抢到锁的 worker 才执行,其余静默跳过。不引入新组件(compose 里没有 redis,redis_client 恒为 None,相关装饰器全程 fail-open,指望不上)。 实测 8 线程并发抢锁恰好 1 个成功。 2) 邮箱授权码被打印进日志 email_service.py 的 print(f"[DEBUG send_email] cfg = {cfg}") 会把含 password 的整个配置打到 stdout → docker logs。改为只打非敏感字段。 新增配置:MAIL_DAILY_REPORT_RECIPIENTS(逗号分隔,默认 duxingchen@iris-rs.cn) --- docker-compose.yml | 2 + .../app/services/daily_report_service.py | 647 ++++++++++++++++++ inventory-backend/app/utils/constants.py | 32 +- inventory-backend/app/utils/email_service.py | 10 +- inventory-backend/app/utils/job_lock.py | 70 ++ inventory-backend/config.py | 4 + inventory-backend/run.py | 57 +- 7 files changed, 812 insertions(+), 10 deletions(-) create mode 100644 inventory-backend/app/services/daily_report_service.py create mode 100644 inventory-backend/app/utils/job_lock.py diff --git a/docker-compose.yml b/docker-compose.yml index 07f17c8..9a3e4f1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -34,6 +34,8 @@ services: MAIL_DEFAULT_SENDER: wms@iris-rs.cn MAIL_USE_SSL: "true" MAIL_USE_TLS: "false" + # MOM 系统日报收件人(逗号分隔可配多个)。每天 17:30 北京时间发送。 + MAIL_DAILY_REPORT_RECIPIENTS: ${MAIL_DAILY_REPORT_RECIPIENTS:-duxingchen@iris-rs.cn} # Track 系统 Webhook 通知(未配置则后端静默跳过,不通知) TRACK_WEBHOOK_URL: ${TRACK_WEBHOOK_URL:-} TRACK_WEBHOOK_KEY: ${TRACK_WEBHOOK_KEY:-} diff --git a/inventory-backend/app/services/daily_report_service.py b/inventory-backend/app/services/daily_report_service.py new file mode 100644 index 0000000..008a625 --- /dev/null +++ b/inventory-backend/app/services/daily_report_service.py @@ -0,0 +1,647 @@ +""" +MOM 系统日报 —— 采集当日数据 → 渲染纯文本 → 发送邮件。 + +设计取舍 +-------- +★ **纯代码统计,不接 AI**。 + + 日报是固定格式的结构化统计(计数 + 清单),不是创造性任务。此前那份日报由 + Dify(外部 AI 服务,见 api/v1/ai_proxy.py)生成,好处是能写人话,代价是: + · 外部依赖一断,整份日报就没了(这正是"AI 坏了"之后发生的事情); + · 按次计费、要走网络、要管密钥; + · **会在报表里编数字** —— 日报是给人做决策用的,数字必须可信。 + + 换成确定性模板后,这些代价全部消失,而日报本来也不需要创造力。 + +★ 口径与前端「操作审计日志」页保持一致,避免两处对同一天给出不同数字: + · action 一律经 canon_action() 归一化(历史数据里 CREATE 与「新增」混用, + 不归一会漏掉早期数据); + · 默认排除 system 占位账号(与审计页默认的"真实用户"视图一致)。 + 注意这是**占位账号**,不是"系统自动产生的日志" —— 判据是 username。 + +★ 凡有上限的地方一律显式写明「另有 N 条」,不静默截断。报表的读法就是 + "看到的就是全部",偷偷砍掉会让人对数据产生错误信心。 + +时间口径 +-------- +audit_logs.created_at 由 beijing_time() 写入,是 naive 北京时间。 +本模块按**北京时间的自然日**切分 [00:00, 次日 00:00)。 +""" +import logging +from collections import Counter, defaultdict +from datetime import datetime, timedelta + +from app.extensions import db +from app.models.audit import AuditLog +from app.models.base import MaterialBase +from app.models.outbound import TransOutbound +from app.models.transaction import TransBorrow +from app.utils.audit_labels import ( + BOOLEAN_FIELDS, + PERSON_NAME_FIELDS, + USER_ID_FIELDS, + bool_label, + canon_action, + enum_label, + field_label, + id_ref_of, + person_name_label, +) +from app.utils.constants import OutboundType + +logger = logging.getLogger(__name__) + +# 系统占位账号:菜单/权限初始化一类的自动写入挂在这个名下(口径同审计页) +SYSTEM_USERNAME = 'system' + +# 未解析出物料信息时的占位 +UNKNOWN = '-' + + +class DailyReportService: + """日报采集与渲染。所有方法均为无状态,便于手动触发与测试。""" + + # 【修改】只详列「变更字段数 >= 该值」的记录 —— 全列会把日报冲成一屏流水 + DETAIL_MIN_FIELDS = 3 + # 【修改】详列条数上限(超出部分显式报数,不静默丢弃) + DETAIL_LIMIT = 30 + # 【新增】模块清单条数上限 + MODULE_LIMIT = 20 + # 【出库】【借库】单据条数上限 + DOC_LIMIT = 50 + # 变更字段的噪声过滤 —— 见 _changes_of() 的说明 + IGNORED_CHANGE_FIELDS = frozenset({ + 'updated_at', + 'product_image', + 'product_image_remark', + 'manual_link', + 'manual_link_remark', + 'purchase_link', + }) + + # ------------------------------------------------------------------ + # 时间边界 + # ------------------------------------------------------------------ + @staticmethod + def _day_bounds(report_date=None): + """返回 [当天 00:00, 次日 00:00)。report_date 为 date 对象,缺省取今天。""" + if report_date is None: + report_date = datetime.now().date() + start = datetime(report_date.year, report_date.month, report_date.day) + return start, start + timedelta(days=1) + + # ------------------------------------------------------------------ + # 物料信息解析(出库/借库明细要显示名称/规格/分类) + # ------------------------------------------------------------------ + @staticmethod + def _load_material_index(pairs): + """ + 批量解析 {(source_table, stock_id)} → {'name','spec','category'}。 + + ★ 批量而非逐条:一张日报可能涉及上百行明细,逐行查库就是 N+1。 + ★ 取不到就跳过(源库存行可能已被物理删除),由渲染层显示占位符 —— + 不因为缺一个物料名就让整份日报发不出去。 + """ + from app.models.inbound.buy import StockBuy + from app.models.inbound.product import StockProduct + from app.models.inbound.semi import StockSemi + + model_map = { + 'stock_buy': StockBuy, + 'stock_semi': StockSemi, + 'stock_product': StockProduct, + } + + by_table = defaultdict(set) + for st, sid in pairs: + st = (st or '').strip() + if st in model_map and sid: + try: + by_table[st].add(int(sid)) + except (TypeError, ValueError): + continue + + index = {} + for table, ids in by_table.items(): + model = model_map[table] + rows = model.query.filter(model.id.in_(ids)).all() + base_ids = {r.base_id for r in rows if getattr(r, 'base_id', None)} + bases = {} + if base_ids: + bases = { + b.id: b + for b in MaterialBase.query.filter(MaterialBase.id.in_(base_ids)).all() + } + for r in rows: + b = bases.get(getattr(r, 'base_id', None)) + if b is None: + continue + # 分类:material_base.category 本身**已经是完整路径** + #(实测 'IRIS/原材料/结构/非标OS'、'LICA/生产配件'), + # company_name 是它的首段、material_type 是它的末段 —— + # 再拼一遍会得到 'LICA/LICA/生产配件/生产配件' 这种重复。 + # 故直接取用,不拼接。 + index[(table, r.id)] = { + 'name': (b.name or '').strip(), + 'spec': (b.spec_model or '').strip(), + 'category': (b.category or '').strip(), + } + return index + + # ------------------------------------------------------------------ + # 采集 + # ------------------------------------------------------------------ + @staticmethod + def _operator_of(row): + """ + 操作人显示名:显示名(账号),与旧日报的「平板(pingban)」同形。 + + ★ display_name 库里存法**不统一**,实测有四种: + '高雪(gaoxue)' / '高闯/gaochuang' / '杜邢宸(duxingchen)' / ''(空) + 直接 f"{display_name}({username})" 会拼出 + '杜邢宸(duxingchen)(duxingchen)' 这种账号重复的怪名字。 + 故先把 display_name 里已内嵌的账号部分剥掉,再统一格式化。 + """ + dn = (getattr(row, 'display_name', '') or '').strip() + un = (getattr(row, 'username', '') or '').strip() + + # 剥掉内嵌账号:'高雪(gaoxue)'→'高雪','高闯/gaochuang'→'高闯' + for sep in ('(', '/'): + if sep in dn: + dn = dn.split(sep)[0].strip() + + if dn and un and dn != un: + return f"{dn}({un})" + return dn or un or UNKNOWN + + @staticmethod + def _fmt_name(raw): + """ + 统一库里五花八门的操作人写法。 + + audit_logs 侧有 display_name/username 两个字段可以拼,但 + trans_outbound.operator_name / trans_borrow.dispatch_operator 这类 + **流水表**只存了一个字符串,实测有 '高闯/gaochuang' 与 + '杜邢宸(duxingchen)' 两种风格。同一封邮件里两种风格并存会显得很随意, + 故统一成「名(账号)」。 + """ + s = (raw or '').strip() + if not s: + return UNKNOWN + if '/' in s: + name, _, acct = s.partition('/') + name, acct = name.strip(), acct.strip() + if name and acct: + return f"{name}({acct})" + return name or acct or s + return s + + @staticmethod + def _as_int(value): + """尽力转 int;转不了返回 None(bool 不算数 —— 它是 int 的子类)""" + if isinstance(value, bool): + return None + if isinstance(value, int): + return value + if isinstance(value, float) and value.is_integer(): + return int(value) + if isinstance(value, str): + s = value.strip() + if s.lstrip('-').isdigit(): + return int(s) + return None + + @staticmethod + def _load_user_map(user_ids): + """{user_id: '名(账号)'}""" + if not user_ids: + return {} + from app.models.system import SysUser + + rows = SysUser.query.filter(SysUser.id.in_(user_ids)).all() + return {u.id: DailyReportService._fmt_name(u.username) for u in rows} + + @staticmethod + def _load_material_name_map(base_ids): + """{base_id: 物料名}""" + if not base_ids: + return {} + rows = MaterialBase.query.filter(MaterialBase.id.in_(base_ids)).all() + return {m.id: (m.name or '').strip() for m in rows} + + @staticmethod + def _load_menu_map(menu_ids): + """ + {menu_id: 菜单名}。 + + ★ 0 是 sys_menu.parent_id 的默认值,含义是「没有上级」= 顶级菜单 —— + 不能当成"查不到的 ID"而显示成 0。 + """ + mapping = {0: '顶级菜单'} + ids = {i for i in menu_ids if i} + if not ids: + return mapping + from app.models.system import SysMenu + + for m in SysMenu.query.filter(SysMenu.id.in_(ids)).all(): + mapping[m.id] = (m.name or '').strip() or f'菜单#{m.id}' + return mapping + + @staticmethod + def _load_return_ledger_map(return_ids): + """ + {trans_return.id: 'SKU 的退回流水(不良品 1.0)'}。 + + 退回流水表本身没有可读单号,能标识它这段记录的就是「哪个物料、什么类型、 + 多少数量」——所以直接把这三样拼出来,而不是显示一个光秃秃的行号。 + """ + if not return_ids: + return {} + from app.models.transaction import TransReturn + + out = {} + for r in TransReturn.query.filter(TransReturn.id.in_(return_ids)).all(): + qty = float(r.return_qty or 0) + rtype = (r.return_type or '').strip() + out[r.id] = f"{r.sku or '-'} 的退回流水({rtype} {qty})".replace('( ', '(') + return out + + @staticmethod + def _resolve_value(module, field, value, ref_maps): + """ + 把单个变更值翻成人类可读形式;翻不动返回 None(调用方回落原值)。 + + 四类映射,优先级即从上到下: + · 布尔列 → 是 / 否 + · 人名串 → '杜邢宸/duxingchen' → '杜邢宸(duxingchen)' + · 引用 ID 的列 → 查实体(用户 / 物料 / 菜单 / 退回流水…, + 指向哪张表由 audit_labels.id_ref_of 按模块判定) + · 枚举列 → 按模块查状态表(见 audit_labels.enum_label) + + ★ 翻不动就**原样显示**,绝不硬编一个可能错的中文 —— 报表里的错译 + 比裸数字更难被发现,也更危险。 + """ + if field in BOOLEAN_FIELDS: + text = bool_label(value) + if text is not None: + return text + + # 人名串('杜邢宸/duxingchen')→ '杜邢宸(duxingchen)'。 + # 按字段名限定,理由见 audit_labels.PERSON_NAME_FIELDS 的说明。 + if field in PERSON_NAME_FIELDS: + text = person_name_label(value) + if text is not None: + return text + + iv = DailyReportService._as_int(value) + if iv is not None: + if field in USER_ID_FIELDS: + return (ref_maps.get('user') or {}).get(iv) + ref = id_ref_of(module, field) + if ref: + got = (ref_maps.get(ref) or {}).get(iv) + if got is not None: + return got + + return enum_label(module, field, value) + + @staticmethod + def _changes_of(row): + """ + UPDATE 记录的字段变更 → [(字段, 旧值, 新值)];结构异常时返回空列表。 + + ★ 噪声字段在**计数之前**就滤掉:图片/链接/更新时间的变化与业务动作无关, + 算进「变更字段数」会把噪声顶过阈值,让日报里塞满 URL 对比 —— + 看起来"有 3 个字段变了",实际业务上什么都没发生。 + """ + details = row.details or {} + if not isinstance(details, dict): + return [] + changes = details.get('changes') + if not isinstance(changes, dict): + return [] + out = [] + for key, val in changes.items(): + if key in DailyReportService.IGNORED_CHANGE_FIELDS: + continue + if isinstance(val, dict): + out.append((key, val.get('old'), val.get('new'))) + else: + # 兼容「非 {old,new} 结构」的历史写法:整值视为新值 + out.append((key, None, val)) + return out + + @staticmethod + def collect(report_date=None): + """ + 采集指定日期的日报数据。 + + Returns: dict(结构见下方各处注释),供 render() 使用,也可单独用于调试。 + """ + start, end = DailyReportService._day_bounds(report_date) + + # ---------- 审计日志 ---------- + audit_rows = AuditLog.query.filter( + AuditLog.created_at >= start, + AuditLog.created_at < end, + AuditLog.username != SYSTEM_USERNAME, + ).all() + + by_action = defaultdict(list) + for r in audit_rows: + by_action[canon_action(r.action)].append(r) + + create_by_module = Counter((r.module or '未知模块') for r in by_action['CREATE']) + + # 【修改】按操作人分组,只保留「变更字段数 >= 阈值」的记录供详列 + update_by_operator = defaultdict(list) + for r in by_action['UPDATE']: + update_by_operator[DailyReportService._operator_of(r)].append(r) + + # ---------- 变更值里的 ID 批量解析 ---------- + # 「实际审批人ID: 空 → 7」这种显示没有意义,要变成「杜邢宸」。 + # 逐条查库是 N+1,故先把涉及的 ID 按类型收集起来,每类一次查完。 + user_ids = set() + ref_ids = defaultdict(set) # ref 类型 -> {id} + for rows in update_by_operator.values(): + for r in rows: + for key, old, new in DailyReportService._changes_of(r): + for v in (old, new): + iv = DailyReportService._as_int(v) + if iv is None: + continue + if key in USER_ID_FIELDS: + user_ids.add(iv) + continue + ref = id_ref_of(r.module, key) + if ref: + ref_ids[ref].add(iv) + + loaders = { + 'user': DailyReportService._load_user_map, + 'material': DailyReportService._load_material_name_map, + 'menu': DailyReportService._load_menu_map, + 'return_ledger': DailyReportService._load_return_ledger_map, + } + ref_maps = {'user': loaders['user'](user_ids)} + for ref, ids in ref_ids.items(): + if ref in loaders: + ref_maps[ref] = loaders[ref](ids) + + # ---------- 出库 ---------- + outbound_rows = TransOutbound.query.filter( + TransOutbound.outbound_time >= start, + TransOutbound.outbound_time < end, + ).order_by(TransOutbound.outbound_time).all() + + # ---------- 借库 ---------- + borrow_rows = TransBorrow.query.filter( + TransBorrow.borrow_time >= start, + TransBorrow.borrow_time < end, + ).order_by(TransBorrow.borrow_time).all() + + # ---------- 物料信息批量解析 ---------- + pairs = [(r.source_table, r.stock_id) for r in outbound_rows] + pairs += [(r.source_table, r.stock_id) for r in borrow_rows] + material_index = DailyReportService._load_material_index(pairs) + + def _item_of(row): + info = material_index.get( + ((row.source_table or '').strip(), row.stock_id), {} + ) + return { + 'sku': row.sku or '', + 'name': info.get('name') or '', + 'spec': info.get('spec') or '', + 'category': info.get('category') or '', + 'quantity': float(row.quantity or 0), + } + + # 出库按单号聚合 + outbound_docs = defaultdict(list) + for r in outbound_rows: + outbound_docs[r.outbound_no or f"#{r.id}"].append(r) + outbound_list = [] + for no, rows in outbound_docs.items(): + head = rows[0] + outbound_list.append({ + 'no': no, + 'operator': DailyReportService._fmt_name(head.operator_name), + 'consumer': head.consumer_name or '', + 'type': OutboundType.label(head.outbound_type), + 'time': head.outbound_time.strftime('%H:%M') if head.outbound_time else '', + 'items': [_item_of(r) for r in rows], + }) + outbound_list.sort(key=lambda d: d['time']) + + # 借库按单号聚合 + borrow_docs = defaultdict(list) + for r in borrow_rows: + borrow_docs[r.borrow_no or f"#{r.id}"].append(r) + borrow_list = [] + for no, rows in borrow_docs.items(): + head = rows[0] + borrow_list.append({ + 'no': no, + 'borrower': DailyReportService._fmt_name(head.borrower_name), + 'operator': DailyReportService._fmt_name(head.dispatch_operator), + 'returned': bool(head.is_returned), + 'time': head.borrow_time.strftime('%H:%M') if head.borrow_time else '', + 'items': [_item_of(r) for r in rows], + }) + borrow_list.sort(key=lambda d: d['time']) + + return { + 'report_date': start.strftime('%Y-%m-%d'), + 'audit': { + 'create_total': len(by_action['CREATE']), + 'update_total': len(by_action['UPDATE']), + 'delete_total': len(by_action['DELETE']), + 'create_by_module': create_by_module.most_common(), + 'update_by_operator': { + k: v for k, v in sorted( + update_by_operator.items(), key=lambda kv: -len(kv[1]) + ) + }, + # 变更值里 ID → 实体的解析结果(key 为 ref 类型),渲染时用 + 'ref_maps': ref_maps, + }, + 'outbound': outbound_list, + 'borrow': borrow_list, + } + + # ------------------------------------------------------------------ + # 渲染 + # ------------------------------------------------------------------ + @staticmethod + def _fmt_value(val): + """变更值显示:空值统一显示「空」(与审计页详情弹窗一致)。""" + if val is None or val == '': + return '空' + s = str(val) + return s if len(s) <= 40 else s[:40] + '…' + + @staticmethod + def _render_changes(lines, module, changes, ref_maps): + """ + 输出一组字段变更:值先经 _resolve_value 翻译,翻不动才显示原值。 + + ★ 「实际审批人ID: 空 → 7」解析成人名后,标签里的「ID」后缀就不贴切了, + 故解析成功时去掉后缀 → 「实际审批人: 空 → 杜邢宸(duxingchen)」。 + """ + for key, old, new in changes: + old_t = DailyReportService._resolve_value(module, key, old, ref_maps) + new_t = DailyReportService._resolve_value(module, key, new, ref_maps) + + label = field_label(key) + is_ref = key in USER_ID_FIELDS or id_ref_of(module, key) is not None + if is_ref and (old_t is not None or new_t is not None) and label.endswith('ID'): + label = label[:-2] + + lines.append( + f" · {label}: " + f"{DailyReportService._fmt_value(old if old_t is None else old_t)} → " + f"{DailyReportService._fmt_value(new if new_t is None else new_t)}" + ) + + @staticmethod + def _render_audit(lines, audit): + # ---------- 新增 ---------- + lines.append("【新增】") + total = audit['create_total'] + lines.append(f" 总量: {total} 条") + if total == 0: + lines.append(" 今日无新增记录") + else: + mods = audit['create_by_module'] + for name, cnt in mods[:DailyReportService.MODULE_LIMIT]: + lines.append(f" · {name}: {cnt} 条") + if len(mods) > DailyReportService.MODULE_LIMIT: + rest = len(mods) - DailyReportService.MODULE_LIMIT + lines.append(f" (另有 {rest} 个模块未列出)") + lines.append("") + + # ---------- 修改 ---------- + lines.append("【修改】") + total = audit['update_total'] + lines.append(f" 总量: {total} 条") + if total == 0: + lines.append(" 今日无修改记录") + lines.append("") + return + # ★ 只详列「有实质变更」的操作人;其余合并成一行报数。 + # 否则每天几十位操作人各占两行「无变更字段≥3的记录」,日报一半是废话。 + quiet_ops = 0 + quiet_records = 0 + for operator, rows in audit['update_by_operator'].items(): + detailed = [r for r in rows + if len(DailyReportService._changes_of(r)) + >= DailyReportService.DETAIL_MIN_FIELDS] + if not detailed: + quiet_ops += 1 + quiet_records += len(rows) + continue + lines.append(f" 【{operator}】共 {len(rows)} 条") + lines.append(f" 以下记录变更字段≥{DailyReportService.DETAIL_MIN_FIELDS}:") + for r in detailed[:DailyReportService.DETAIL_LIMIT]: + lines.append(f" · [{r.module or '未知模块'}] {r.url or ''}") + DailyReportService._render_changes( + lines, r.module, DailyReportService._changes_of(r), + audit.get('ref_maps') or {}, + ) + if len(detailed) > DailyReportService.DETAIL_LIMIT: + rest = len(detailed) - DailyReportService.DETAIL_LIMIT + lines.append(f" (另有 {rest} 条未列出)") + lines.append("") + + if quiet_ops: + lines.append( + f" 另有 {quiet_ops} 位操作人(共 {quiet_records} 条)" + f"无变更字段≥{DailyReportService.DETAIL_MIN_FIELDS}的记录,未逐条列出" + ) + lines.append("") + + @staticmethod + def _render_docs(lines, title, docs, head_fields, empty_text): + lines.append(f"【{title}】") + lines.append(f" 总量: {len(docs)} 条") + if not docs: + lines.append(f" {empty_text}") + lines.append("") + return + for d in docs[:DailyReportService.DOC_LIMIT]: + lines.append(f" · {d['no']}") + lines.append(" " + " | ".join(f"{k}: {d[v]}" for k, v in head_fields)) + for it in d['items']: + lines.append( + f" 物料: {it['name']} | 规格: {it['spec'] or UNKNOWN} | " + f"数量: {it['quantity']} | 分类: {it['category'] or UNKNOWN}" + ) + if len(docs) > DailyReportService.DOC_LIMIT: + rest = len(docs) - DailyReportService.DOC_LIMIT + lines.append(f" (另有 {rest} 张单据未列出)") + lines.append("") + + @staticmethod + def render(stats): + """把 collect() 的结果渲染成纯文本正文。返回 (subject, body)。""" + date_str = stats['report_date'] + lines = [f"📋 MOM系统日报 — {date_str}", "=" * 60, ""] + + DailyReportService._render_audit(lines, stats['audit']) + + lines.append("【删除】") + lines.append(f" 总量: {stats['audit']['delete_total']} 条") + lines.append("") + + DailyReportService._render_docs( + lines, "出库", stats['outbound'], + head_fields=[('操作员', 'operator'), ('客户', 'consumer'), ('类型', 'type')], + empty_text="今日无出库记录", + ) + DailyReportService._render_docs( + lines, "借库", stats['borrow'], + head_fields=[('借用人', 'borrower'), ('库管', 'operator')], + empty_text="今日无借库记录", + ) + + lines.append("=" * 60) + lines.append("此邮件由 MOM 系统自动发送,请勿回复。") + + subject = f"MOM系统日报 {date_str}" + return subject, '\n'.join(lines) + + # ------------------------------------------------------------------ + # 发送 + # ------------------------------------------------------------------ + @staticmethod + def build_report(report_date=None): + """采集 + 渲染,但不发送 —— 用于预览与调试。返回 (subject, body, stats)。""" + stats = DailyReportService.collect(report_date) + subject, body = DailyReportService.render(stats) + return subject, body, stats + + @staticmethod + def send_daily_report(report_date=None, recipients=None): + """ + 采集 → 渲染 → 发送。供定时任务与手动触发共用。 + + recipients 缺省时取 config.MAIL_DAILY_REPORT_RECIPIENTS。 + 返回 {'subject','recipients','stats'},便于调用方打日志。 + """ + from flask import current_app + from app.utils.email_service import send_email + + if recipients is None: + raw = current_app.config.get('MAIL_DAILY_REPORT_RECIPIENTS', '') or '' + recipients = [e.strip() for e in raw.split(',') if e.strip()] + if isinstance(recipients, str): + recipients = [e.strip() for e in recipients.split(',') if e.strip()] + if not recipients: + logger.warning("[日报] 未配置收件人(MAIL_DAILY_REPORT_RECIPIENTS),跳过发送") + return {'subject': None, 'recipients': [], 'stats': None} + + subject, body, stats = DailyReportService.build_report(report_date) + send_email(recipients, subject, body) + logger.info(f"[日报] 已发送 {stats['report_date']} → {recipients}") + return {'subject': subject, 'recipients': recipients, 'stats': stats} diff --git a/inventory-backend/app/utils/constants.py b/inventory-backend/app/utils/constants.py index 3b21c08..457999a 100644 --- a/inventory-backend/app/utils/constants.py +++ b/inventory-backend/app/utils/constants.py @@ -24,4 +24,34 @@ class UserRole: OUTBOUND: '出库员', PURCHASER: '采购员', SALES: '销售' - } \ No newline at end of file + } + + +class OutboundType: + """ + 出库类型(trans_outbound.outbound_type)。 + + ★ 前端 views/outbound/index.vue 的 formatType() 里另有一份同名映射, + 两份是手工同步的副本。日报(后端)需要同一份口径,故在此定义; + **改这里时请一并改前端**,或后续像审计标签那样改为接口下发。 + """ + SALES = 'SALES' # 销售出库 + USE = 'USE' # 内部领用 + PRODUCTION = 'PRODUCTION' # 生产出库 + SCRAP = 'SCRAP' # 报废 + LOSS = 'LOSS' # 盘亏出库 + REPAIR = 'REPAIR' # 维修出库 + + LABELS = { + SALES: '销售出库', + USE: '内部领用', + PRODUCTION: '生产出库', + SCRAP: '报废', + LOSS: '盘亏出库', + REPAIR: '维修出库', + } + + @classmethod + def label(cls, code): + """未命中时原样返回,避免报表上出现空白""" + return cls.LABELS.get((code or '').strip(), (code or '').strip()) \ No newline at end of file diff --git a/inventory-backend/app/utils/email_service.py b/inventory-backend/app/utils/email_service.py index 0f20839..3923180 100644 --- a/inventory-backend/app/utils/email_service.py +++ b/inventory-backend/app/utils/email_service.py @@ -60,7 +60,15 @@ def send_email(to_email: Union[str, List[str]], subject: str, content: str, cfg: if cfg is None: cfg = _get_config() - print(f"[DEBUG send_email] cfg = {cfg}") + # ★ 不要把 cfg 整个打出来 —— 它含 'password'(邮箱授权码), + # 一 print 就落到 stdout → docker logs,任何能看日志的人都拿得到发信凭证。 + # 只打非敏感字段,密码只报"有没有设"。 + _pw_state = '已设置' if cfg.get('password') else '空' + print( + f"[DEBUG send_email] server={cfg.get('server')} port={cfg.get('port')} " + f"sender={cfg.get('sender')} ssl={cfg.get('use_ssl')} tls={cfg.get('use_tls')} " + f"enabled={cfg.get('enabled')} password={_pw_state}" + ) # 发送总开关 if not cfg.get('enabled'): diff --git a/inventory-backend/app/utils/job_lock.py b/inventory-backend/app/utils/job_lock.py new file mode 100644 index 0000000..d731d19 --- /dev/null +++ b/inventory-backend/app/utils/job_lock.py @@ -0,0 +1,70 @@ +""" +跨进程的定时任务互斥锁。 + +为什么需要这个 +-------------- +`gunicorn.conf.py` 配了 8 个 worker(`workers = min(cpu*2+1, 8)`),而 `run.py` +是在**模块级**启动 APScheduler 的 —— gunicorn 未开 preload_app,每个 worker +都会独立 import 一次 `run.py`,于是**每个 worker 各起一份调度器**,同一个 +cron 任务在相同时刻被并发执行 8 次。 + +实测:容器里跑着 8 个 worker 进程,`create_app()` 与"调度器已启动"两条日志 +的出现次数完全相等 —— 每个 worker 都配了一份。 + +后果:库存预警邮件每天实际会重复发 8 封。 + +APScheduler 自带的 `max_instances` 只在**单个调度器实例内**生效,管不了 +跨进程;`replace_existing` 同理。这里用 PostgreSQL 的**会话级咨询锁** +(`pg_try_advisory_lock`)做真正的跨进程互斥:抢到锁的那个 worker 才执行, +其余直接跳过。好处是不用引入 Redis 之类的新组件(compose 里也没有 redis, +`redis_client` 恒为 None,相关装饰器全程 fail-open,指望不上)。 + +★ 会话级锁绑定在**连接**上,必须在同一条连接上解锁 —— 所以这里显式取一条 + 专用连接,用完归还,不复用请求上下文里的 session。 +""" +import logging +from contextlib import contextmanager + +from sqlalchemy import text + +from app.extensions import db + +logger = logging.getLogger(__name__) + +# 锁标识:PostgreSQL 咨询锁是 bigint,取固定值便于排查(不要用随机数)。 +LOCK_INVENTORY_WARNING = 891001001 # 库存预警每日邮件 +LOCK_DAILY_REPORT = 891001002 # MOM 系统日报 + + +@contextmanager +def advisory_lock(key): + """ + 尝试获取会话级咨询锁,产出「是否抢到」。 + + 用法:: + + with advisory_lock(LOCK_DAILY_REPORT) as acquired: + if not acquired: + return # 别的 worker 正在跑,本轮跳过 + ...实际干活... + + ★ 抢不到锁是**正常路径**,不是错误 —— 说明另一个 worker 正在执行同一任务。 + 调用方应当静默跳过,而不是报警。 + ★ 拿锁失败(数据库不可用等)会让异常向上抛:宁可让调度器记一次失败, + 也不要"没抢到锁"和"抢锁时数据库挂了"两种截然不同的情况被混为一谈。 + """ + conn = db.engine.connect() + acquired = False + try: + acquired = bool( + conn.execute(text("SELECT pg_try_advisory_lock(:k)"), {"k": key}).scalar() + ) + yield acquired + finally: + if acquired: + try: + conn.execute(text("SELECT pg_advisory_unlock(:k)"), {"k": key}) + except Exception as e: # noqa: BLE001 + # 解锁失败不应盖过业务异常;连接关闭时锁也会随之释放 + logger.warning(f"[JobLock] 释放咨询锁 {key} 失败: {e}") + conn.close() diff --git a/inventory-backend/config.py b/inventory-backend/config.py index 7c4e5c5..53ae24c 100644 --- a/inventory-backend/config.py +++ b/inventory-backend/config.py @@ -70,6 +70,10 @@ class Config: MAIL_DEFAULT_SENDER = os.getenv('MAIL_DEFAULT_SENDER', 'wms@iris-rs.cn') # 是否启用邮件发送功能(开发环境可设为 false 禁用) MAIL_ENABLED = os.getenv('MAIL_ENABLED', 'true').lower() in ('true', '1', 'yes') + # MOM 系统日报收件人(逗号分隔多个)。留空则日报任务跳过发送并在日志告警。 + MAIL_DAILY_REPORT_RECIPIENTS = os.getenv( + 'MAIL_DAILY_REPORT_RECIPIENTS', 'duxingchen@iris-rs.cn' + ) # ========================================================= # 7. Track 系统 Webhook 通知配置 (发送侧) diff --git a/inventory-backend/run.py b/inventory-backend/run.py index f0960d0..88d441b 100644 --- a/inventory-backend/run.py +++ b/inventory-backend/run.py @@ -4,7 +4,18 @@ from app import create_app app = create_app() # ========================================================= -# 启动时注册库存预警定时任务(每天 9:30 北京时) +# 定时任务注册 +# +# ★★ 每个任务都必须套 advisory_lock,否则会被执行 **8 次**。 +# +# gunicorn.conf.py 配了 8 个 worker,且未开 preload_app —— 每个 worker +# 都会独立 import 一次本文件,于是在模块级启动的这份调度器会**每个 +# worker 各起一份**,同一个 cron 任务在相同时刻被并发执行 8 次。 +# (库存预警此前正是如此:5 条预警配置,每天实际发 8 封重复邮件。 +# APScheduler 的 max_instances 只管单个调度器实例内,管不了跨进程。) +# +# advisory_lock 用 PostgreSQL 会话级咨询锁做真正的跨进程互斥, +# 抢到锁的 worker 才执行,其余静默跳过。详见 app/utils/job_lock.py。 # ========================================================= from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger @@ -14,13 +25,36 @@ beijing_tz = pytz.timezone('Asia/Shanghai') def _run_warning_job(): + """库存预警扫描与邮件发送(每天 9:30 北京时间)""" with app.app_context(): - try: - from app.services.inventory_task import InventoryWarningService - result = InventoryWarningService.check_and_send_warning_emails() - print(f"[Scheduler] 库存预警扫描完成: red={result['red_count']}, yellow={result['yellow_count']}") - except Exception as e: - print(f"[Scheduler] 库存预警任务失败: {e}") + from app.utils.job_lock import advisory_lock, LOCK_INVENTORY_WARNING + with advisory_lock(LOCK_INVENTORY_WARNING) as acquired: + if not acquired: + # 正常路径:另一个 worker 正在执行同一任务 + print("[Scheduler] 库存预警:另一 worker 正在执行,本轮跳过") + return + try: + from app.services.inventory_task import InventoryWarningService + result = InventoryWarningService.check_and_send_warning_emails() + print(f"[Scheduler] 库存预警扫描完成: red={result['red_count']}, yellow={result['yellow_count']}") + except Exception as e: + print(f"[Scheduler] 库存预警任务失败: {e}") + + +def _run_daily_report_job(): + """MOM 系统日报(每天 17:30 北京时间)""" + with app.app_context(): + from app.utils.job_lock import advisory_lock, LOCK_DAILY_REPORT + with advisory_lock(LOCK_DAILY_REPORT) as acquired: + if not acquired: + print("[Scheduler] 系统日报:另一 worker 正在执行,本轮跳过") + return + try: + from app.services.daily_report_service import DailyReportService + r = DailyReportService.send_daily_report() + print(f"[Scheduler] 系统日报已发送: {r['subject']} -> {r['recipients']}") + except Exception as e: + print(f"[Scheduler] 系统日报任务失败: {e}") scheduler = BackgroundScheduler(timezone=beijing_tz) @@ -31,8 +65,15 @@ scheduler.add_job( name='库存预警每日邮件发送', replace_existing=True ) +scheduler.add_job( + func=_run_daily_report_job, + trigger=CronTrigger(hour=17, minute=30, timezone=beijing_tz), + id='daily_report', + name='MOM系统日报', + replace_existing=True +) scheduler.start() -print("✅ 库存预警定时任务已启动(每天 9:30 北京时间执行)") +print("✅ 定时任务已启动:库存预警 9:30 / MOM系统日报 17:30(北京时间)") if __name__ == '__main__': # =================================================