Files
KCGL/inventory-backend/app/services/outbound_service.py
yueli b57c21a4cd feat(inventory): 库存预占生命周期,消除出库/借库超卖
问题:库存超卖
--------------
改造前出库/借库申请只记录「要什么、要多少」,不绑定具体库存行,
真正的 available_quantity 扣减发生在执行阶段。于是多张申请可以同时
claim 同一批货,等到工人拿扫码枪时才发现货已被别人领走。

生命周期(三阶段)
------------------
  提交申请(预占)  reserve_for_items()
     用分配器把需求落到具体库存行,立即扣减 available_quantity,
     并把 (stock_id, source_table, allocated_qty, reserved) 写回 items_json。

  驳回(释放)      release_reserved()
     遍历 items_json 把预占量还回池子,避免货被永不执行的单永久占住。

  扫码执行(覆盖)  verify_scanned() + restore_then_deduct()
     校验实扫身份/数量未超批准范围 → 释放全部预占 → 对实扫批次
     同时扣减 available_quantity 与 stock_quantity。

身份键:base_id 主键 + SKU 兜底(重要设计决策)
-----------------------------------------------
本系统中 SKU 是**批次级**编号:同一 base_id 下每个入库批次各有不同的
SKU(实测 stock_buy 有 183 个物料是多批次的,如 base_id=2405 下有
0000001685 与 0000001974 两个 SKU)。

若以 SKU 作为身份主键,「申请时锁定 A 批、工人现场改扫 B 批」会被判为
身份不符而拒绝 —— 恰好否定了「物理覆盖」这个核心能力。
故改用 base_id(物料级、跨批次稳定,spec_model 由其唯一确定),
历史数据无 base_id 时降级为 (name, spec_model)。

可用量校验按物料汇总,而非按单批次
----------------------------------
开发中修正的一处缺陷:若逐行要求「该批次可用量 >= 该批次扫码量」,
工人改扫小批次时会被误拒。例如本单预占 A 批 5 件,改扫 B 批 2 件 +
C 批 3 件,B 批自身只有 2 件可用,逐行校验即失败。实际这 5 件都是本单
锁定的货,理应允许。现按物料汇总校验可用量,按行校验实物库存。

改动文件
--------
· 新增 app/services/inventory_reservation.py(通用服务层)
· outbound_service.create_request     —— Phase 1 预占
· outbound_service.approve(reject)    —— Phase 2 释放
· outbound_service.create_outbound_batch —— Phase 3 覆盖(移除原逐行扣减)
· borrow_service.submit_approval      —— Phase 1
· borrow_service.approve(reject)      —— Phase 2
· trans_service.execute_dispatch      —— Phase 3,并用统一身份键替换
                                         原有的 (name, spec_model) 字符串匹配

实测(真实 HTTP 全链路)
------------------------
初始 available=10
  ① 提交申请(需5)   → 200,available 10→5    预占生效
  ② 审批通过        → available 仍为 5        预占保留
  ③ 扫码执行(改扫另一批次 4 件) → 200
     原批次恢复满额、实扫批次扣减(0,0),available=6, stock=6

单场景验证:预占 A 批改扫 B 批放行;驳回后可用量完全恢复;
           扫其他物料被拒;批准 6 扫 8 被拒。
2026-09-10 14:59:10 +08:00

1311 lines
55 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.

import uuid # .material -> .base refactor checked
from datetime import datetime, timezone, timedelta
from sqlalchemy import or_, func, desc, and_
from sqlalchemy.orm import joinedload
from app.extensions import db, beijing_time
from app.models.outbound import TransOutbound, OutboundApproval
# 引入所有库存模型以进行查询
from app.models.inbound.buy import StockBuy
from app.models.inbound.semi import StockSemi
from app.models.inbound.product import StockProduct
# 引入基础信息表
from app.models.base import MaterialBase
# 引入维修单表
from app.models.transaction import TransRepair
# 引入系统用户表
from app.models.system import SysUser
# Track 联动通知(出库 → Track 标记"已出库")
from app.services.track_webhook_service import notify_track, get_current_operator
class OutboundService:
@staticmethod
def generate_outbound_no():
"""
生成出库单号: OUT-yyyyMMdd-HHmm-当日流水(4位)
例如: OUT-20260205-1558-0001
"""
beijing_tz = timezone(timedelta(hours=8))
now = datetime.now(beijing_tz)
date_str = now.strftime('%Y%m%d')
time_str = now.strftime('%H%M')
prefix = f"OUT-{date_str}-"
existing_count = db.session.query(func.count(func.distinct(TransOutbound.outbound_no))) \
.filter(TransOutbound.outbound_no.like(f"{prefix}%")).scalar()
sequence = existing_count + 1
return f"OUT-{date_str}-{time_str}-{sequence:04d}"
@staticmethod
def get_stock_by_barcode(barcode):
"""
根据扫码内容查找对应的库存物品,并附带价格信息
"""
if not barcode:
return None
clean_code = barcode.strip()
def get_price(item, table_type):
if table_type == 'stock_product':
return float(item.sale_price) if item.sale_price else 0
elif table_type == 'stock_buy':
return float(item.pre_tax_unit_price) if item.pre_tax_unit_price else 0
return 0
# ★ 行级公司隔离:扫码出库/借库只能命中本公司的库存(超管/跨域不受限)
from app.utils.decorators import get_current_company_filter
from app.models.base import MaterialBase
company_limit = get_current_company_filter()
prod_q = StockProduct.query.filter(
or_(StockProduct.barcode == clean_code, StockProduct.sku == clean_code)
)
if company_limit is not None:
prod_q = prod_q.join(MaterialBase, StockProduct.base_id == MaterialBase.id) \
.filter(MaterialBase.company_name == company_limit)
prod = prod_q.first()
if prod:
res = OutboundService._format_scan_result(prod, 'stock_product')
res['price'] = get_price(prod, 'stock_product')
return res
semi_q = StockSemi.query.filter(
or_(StockSemi.barcode == clean_code, StockSemi.sku == clean_code)
)
if company_limit is not None:
semi_q = semi_q.join(MaterialBase, StockSemi.base_id == MaterialBase.id) \
.filter(MaterialBase.company_name == company_limit)
semi = semi_q.first()
if semi:
res = OutboundService._format_scan_result(semi, 'stock_semi')
res['price'] = 0
return res
buy_q = StockBuy.query.filter(
or_(StockBuy.barcode == clean_code, StockBuy.sku == clean_code)
)
if company_limit is not None:
buy_q = buy_q.join(MaterialBase, StockBuy.base_id == MaterialBase.id) \
.filter(MaterialBase.company_name == company_limit)
buy = buy_q.first()
if buy:
res = OutboundService._format_scan_result(buy, 'stock_buy')
res['price'] = get_price(buy, 'stock_buy')
return res
# 查询维修单表 (按SKU或序列号查询,排除已出库状态)
repair = TransRepair.query.filter(
or_(TransRepair.sku == clean_code, TransRepair.serial_number == clean_code)
).filter(
TransRepair.repair_status != '已出库'
).first()
if repair:
res = {
'id': repair.id,
'sku': repair.sku,
'name': repair.material_name or "维修件",
'spec_model': "",
'category': "",
'material_type': "",
'source_table': 'trans_repair',
'stock_quantity': 1,
'available_quantity': 1,
'batch_number': repair.serial_number or '',
'serial_number': repair.serial_number or '',
'warehouse_location': repair.customer_location or '',
'barcode': repair.sku,
'price': float(repair.sale_price) if repair.sale_price else 0
}
return res
return None
@staticmethod
def _format_scan_result(item, table_name):
base_name = ""
base_spec = ""
base_cat = ""
base_type = ""
if hasattr(item, 'base') and item.base:
base_name = item.base.name
base_spec = item.base.spec_model
base_cat = item.base.category
base_type = item.base.material_type
if not base_name and hasattr(item, 'base_id') and item.base_id:
try:
base_info = MaterialBase.query.get(item.base_id)
if base_info:
base_name = base_info.name
base_spec = base_info.spec_model
base_cat = base_info.category
base_type = base_info.material_type
except Exception:
pass
if not base_name and hasattr(item, 'base') and item.base:
base_name = item.base.name
stock_qty = float(item.stock_quantity) if item.stock_quantity else 0
avail_qty = float(item.available_quantity) if item.available_quantity else 0
return {
'id': item.id,
'sku': item.sku,
'name': base_name or "未知物品",
'spec_model': base_spec or "",
'category': base_cat or "",
'material_type': base_type or "",
'source_table': table_name,
'stock_quantity': stock_qty,
'available_quantity': avail_qty,
'batch_number': getattr(item, 'batch_number', ''),
'warehouse_location': getattr(item, 'warehouse_location', ''),
'barcode': getattr(item, 'barcode', '')
}
@staticmethod
def create_outbound_batch(data, operator_name='System'):
items = data.get('items', [])
if not items:
raise ValueError("出库商品列表不能为空")
outbound_no = OutboundService.generate_outbound_no()
common_data = {
'outbound_no': outbound_no,
'consumer_name': data.get('consumer_name'),
'outbound_type': data.get('outbound_type', 'SALES'),
'signature_path': data.get('signature_path'),
'operator_name': operator_name,
'remark': data.get('remark')
}
beijing_tz = timezone(timedelta(hours=8))
current_time = datetime.now(beijing_tz).replace(tzinfo=None)
# ★ 审批单相关逻辑
request_id = data.get('request_id')
approval = None
if request_id:
# 根据 request_id 查询审批单
approval = OutboundApproval.query.get(request_id)
if not approval:
raise ValueError(f"关联的审批单不存在 (ID: {request_id})")
if approval.status != 1:
status_map = {0: '待审批', 1: '已通过', 2: '已驳回', 3: '已完成'}
current_status = status_map.get(approval.status, str(approval.status))
raise ValueError(
f"关联的审批单状态不允许出库 (当前状态: {current_status}),"
f"仅已通过的审批单方可执行出库"
)
model_map = {
'stock_buy': StockBuy,
'stock_semi': StockSemi,
'stock_product': StockProduct
}
# ★ Track 联动收集:(serial_number, source_table, quantity)
track_notifications = []
# ==================================================================
# ★ Phase 3:预占再平衡(仅针对关联审批单的出库)
#
# 关联审批单时,申请阶段已把货预占在「申请时选定的批次」上。
# 工人实际扫的可能是同物料的**另一个批次**(物理覆盖),因此这里:
# 1. 校验实扫的身份/数量未超出批准范围(base_id 主键,允许换批次)
# 2. 释放全部预占
# 3. 对实扫批次扣减 available_quantity 与 stock_quantity
# 之后主循环只写 TransOutbound 流水,不再重复扣库存。
#
# 无关联审批单(散单)时跳过,走原有逐行扣减逻辑。
# ==================================================================
if approval is not None:
from app.services.inventory_reservation import (
verify_scanned, restore_then_deduct,
)
_approved = approval.get_items()
# 仅处理库存类来源;维修单(trans_repair)不走库存预占
_scanned = [i for i in items if i.get('source_table') != 'trans_repair']
if _scanned:
verify_scanned(_scanned, _approved)
restore_then_deduct(_scanned, _approved)
try:
for item in items:
source_table = item.get('source_table')
stock_id = item.get('stock_id')
quantity = float(item.get('quantity', 0))
unit_price = float(item.get('price', 0))
if quantity <= 0:
raise ValueError(f"SKU {item.get('sku')} 的出库数量必须大于0")
# 处理维修单出库
if source_table == 'trans_repair':
repair = TransRepair.query.with_for_update().get(stock_id)
if not repair:
raise ValueError(f"维修单不存在 (ID: {stock_id})")
# 更新维修单状态为已出库
repair.repair_status = '已出库'
repair.shipping_date = current_time
# 收集 Track 联动信息(维修单带 serial_number)
track_notifications.append((getattr(repair, 'serial_number', None), source_table, quantity))
# 创建出库记录
new_record = TransOutbound(
sku=item.get('sku'),
source_table=source_table,
stock_id=stock_id,
barcode=item.get('barcode'),
quantity=quantity,
unit_price=unit_price,
outbound_time=current_time,
**common_data
)
db.session.add(new_record)
continue
ModelClass = model_map.get(source_table)
if not ModelClass:
continue
# ==========================================================
# ★ Phase 3:库存扣减已由「预占 + 再平衡」统一处理
#
# 流程(在下方 _apply_reservation_override 中完成):
# 1. 校验实扫身份/数量落在批准范围内(base_id 主键匹配,允许换批次)
# 2. 释放申请时锁定的全部批次(available_quantity 还回池子)
# 3. 对实扫批次扣减 available_quantity 与 stock_quantity
#
# 因此此处**不再**直接扣减库存,避免与再平衡逻辑重复扣两次。
# ==========================================================
stock_record = ModelClass.query.get(stock_id)
if not stock_record:
raise ValueError(f"库存记录不存在 (ID: {stock_id})")
# 收集 Track 联动信息(库存表 serial_number = Track 身份证)
track_notifications.append((getattr(stock_record, 'serial_number', None), source_table, quantity))
new_record = TransOutbound(
sku=item.get('sku'),
source_table=source_table,
stock_id=stock_id,
barcode=item.get('barcode'),
quantity=quantity,
unit_price=unit_price,
outbound_time=current_time,
# [新增] 记录出库时的库位快照
warehouse_location=getattr(stock_record, 'warehouse_location', None),
**common_data
)
db.session.add(new_record)
# ★ 如果关联了审批单,出库成功后更新审批单状态为"已完成"
if approval:
approval.status = 3 # 3-已完成
# updated_at 会在 commit 时由 SQLAlchemy 自动更新
# ★ 先提交事务,释放所有行锁,避免 SMTP 调用延长锁持有时间
db.session.commit()
# ★ 出库后通知 Track(发货出库 → Track 标记"已出库")
# 仅在配置 TRACK_OUTBOUND_WEBHOOK_URL 时生效;notify_track 自身容错不阻断业务
try:
from flask import current_app
outbound_url = current_app.config.get('TRACK_OUTBOUND_WEBHOOK_URL')
for sn, src_tbl, qty in track_notifications:
if not sn:
continue
notify_track({
'event': 'outbound.created',
'source_table': src_tbl,
'serial_number': str(sn),
'quantity': float(qty or 0),
'outbound_type': common_data['outbound_type'],
'operator': get_current_operator(),
}, url=outbound_url)
except Exception as e:
import logging
logging.getLogger(__name__).warning(f"⚠️ Track 出库通知失败: {e}")
# ★ 出库后检查低库存预警(移到 commit 之后,避免事务内网络调用)
try:
from app.services.inventory_task import InventoryWarningService
InventoryWarningService.check_and_send_warning_emails()
except Exception as e:
import logging
logging.getLogger(__name__).warning(f"⚠️ 低库存预警检查失败: {e}")
return outbound_no
except Exception as e:
db.session.rollback()
raise e
@staticmethod
def get_grouped_list(page=1, per_page=10, keyword=None, search_type='all', start_date=None, end_date=None, company=None, consumer_name=None, advanced_filters=None):
"""
查询出库记录(按出库单号分组),包含详细物品信息
支持跨表搜索:单号、领用人、SKU、物料名称、规格型号
search_type: all, no, name, sku, material_name, spec_model
company: 可选的公司过滤参数
advanced_filters: 高级筛选条件列表 [{'field','operator','value'}, ...]
"""
# 日期补全:解决零点截断问题
if end_date and len(str(end_date).strip()) == 10:
end_date = f"{str(end_date).strip()} 23:59:59"
if start_date and len(str(start_date).strip()) == 10:
start_date = f"{str(start_date).strip()} 00:00:00"
# 1. 构建基础查询
# 如果有关键词,需要联表搜索物料名称和规格型号
if keyword:
# 根据 search_type 构建不同的搜索条件
if search_type == 'all':
# 原有逻辑:or_ 联表全局模糊搜索
# 查询 stock_buy 路径匹配的名称/规格
buy_match = db.session.query(TransOutbound.outbound_no).join(
StockBuy, and_(
TransOutbound.stock_id == StockBuy.id,
TransOutbound.source_table == 'stock_buy'
)
).join(
MaterialBase, StockBuy.base_id == MaterialBase.id
).filter(
or_(
MaterialBase.name.ilike(f'%{keyword}%'),
MaterialBase.spec_model.ilike(f'%{keyword}%')
)
).subquery()
# 查询 stock_semi 路径匹配的名称/规格
semi_match = db.session.query(TransOutbound.outbound_no).join(
StockSemi, and_(
TransOutbound.stock_id == StockSemi.id,
TransOutbound.source_table == 'stock_semi'
)
).join(
MaterialBase, StockSemi.base_id == MaterialBase.id
).filter(
or_(
MaterialBase.name.ilike(f'%{keyword}%'),
MaterialBase.spec_model.ilike(f'%{keyword}%')
)
).subquery()
# 查询 stock_product 路径匹配的名称/规格
product_match = db.session.query(TransOutbound.outbound_no).join(
StockProduct, and_(
TransOutbound.stock_id == StockProduct.id,
TransOutbound.source_table == 'stock_product'
)
).join(
MaterialBase, StockProduct.base_id == MaterialBase.id
).filter(
or_(
MaterialBase.name.ilike(f'%{keyword}%'),
MaterialBase.spec_model.ilike(f'%{keyword}%')
)
).subquery()
# 合并三种来源的匹配单号
all_matches = db.session.query(buy_match.c.outbound_no).union(
db.session.query(semi_match.c.outbound_no),
db.session.query(product_match.c.outbound_no)
).subquery()
keyword_conditions = or_(
TransOutbound.outbound_no.ilike(f'%{keyword}%'),
TransOutbound.consumer_name.ilike(f'%{keyword}%'),
TransOutbound.sku.ilike(f'%{keyword}%'),
TransOutbound.outbound_no.in_(all_matches)
)
elif search_type == 'no':
keyword_conditions = TransOutbound.outbound_no.ilike(f'%{keyword}%')
elif search_type == 'name':
keyword_conditions = TransOutbound.consumer_name.ilike(f'%{keyword}%')
elif search_type == 'sku':
keyword_conditions = TransOutbound.sku.ilike(f'%{keyword}%')
elif search_type == 'material_name':
# 联表查询物料名称
buy_match = db.session.query(TransOutbound.outbound_no).join(
StockBuy, and_(
TransOutbound.stock_id == StockBuy.id,
TransOutbound.source_table == 'stock_buy'
)
).join(
MaterialBase, StockBuy.base_id == MaterialBase.id
).filter(MaterialBase.name.ilike(f'%{keyword}%')).subquery()
semi_match = db.session.query(TransOutbound.outbound_no).join(
StockSemi, and_(
TransOutbound.stock_id == StockSemi.id,
TransOutbound.source_table == 'stock_semi'
)
).join(
MaterialBase, StockSemi.base_id == MaterialBase.id
).filter(MaterialBase.name.ilike(f'%{keyword}%')).subquery()
product_match = db.session.query(TransOutbound.outbound_no).join(
StockProduct, and_(
TransOutbound.stock_id == StockProduct.id,
TransOutbound.source_table == 'stock_product'
)
).join(
MaterialBase, StockProduct.base_id == MaterialBase.id
).filter(MaterialBase.name.ilike(f'%{keyword}%')).subquery()
all_matches = db.session.query(buy_match.c.outbound_no).union(
db.session.query(semi_match.c.outbound_no),
db.session.query(product_match.c.outbound_no)
).subquery()
keyword_conditions = TransOutbound.outbound_no.in_(all_matches)
elif search_type == 'spec_model':
# 联表查询规格型号
buy_match = db.session.query(TransOutbound.outbound_no).join(
StockBuy, and_(
TransOutbound.stock_id == StockBuy.id,
TransOutbound.source_table == 'stock_buy'
)
).join(
MaterialBase, StockBuy.base_id == MaterialBase.id
).filter(MaterialBase.spec_model.ilike(f'%{keyword}%')).subquery()
semi_match = db.session.query(TransOutbound.outbound_no).join(
StockSemi, and_(
TransOutbound.stock_id == StockSemi.id,
TransOutbound.source_table == 'stock_semi'
)
).join(
MaterialBase, StockSemi.base_id == MaterialBase.id
).filter(MaterialBase.spec_model.ilike(f'%{keyword}%')).subquery()
product_match = db.session.query(TransOutbound.outbound_no).join(
StockProduct, and_(
TransOutbound.stock_id == StockProduct.id,
TransOutbound.source_table == 'stock_product'
)
).join(
MaterialBase, StockProduct.base_id == MaterialBase.id
).filter(MaterialBase.spec_model.ilike(f'%{keyword}%')).subquery()
all_matches = db.session.query(buy_match.c.outbound_no).union(
db.session.query(semi_match.c.outbound_no),
db.session.query(product_match.c.outbound_no)
).subquery()
keyword_conditions = TransOutbound.outbound_no.in_(all_matches)
else:
keyword_conditions = None
else:
keyword_conditions = None
# 【行级数据隔离】基于 JWT 多租户公司过滤
# 通过三个库存表路径,找到匹配公司的出库单号(排除 trans_repair,因其无 MaterialBase 关联)
from app.utils.decorators import get_current_company_filter
company_limit = get_current_company_filter()
if company_limit is not None:
buy_comp = db.session.query(TransOutbound.outbound_no).join(
StockBuy, and_(
TransOutbound.stock_id == StockBuy.id,
TransOutbound.source_table == 'stock_buy'
)
).join(MaterialBase, StockBuy.base_id == MaterialBase.id).filter(
MaterialBase.company_name == company_limit
).subquery()
semi_comp = db.session.query(TransOutbound.outbound_no).join(
StockSemi, and_(
TransOutbound.stock_id == StockSemi.id,
TransOutbound.source_table == 'stock_semi'
)
).join(MaterialBase, StockSemi.base_id == MaterialBase.id).filter(
MaterialBase.company_name == company_limit
).subquery()
prod_comp = db.session.query(TransOutbound.outbound_no).join(
StockProduct, and_(
TransOutbound.stock_id == StockProduct.id,
TransOutbound.source_table == 'stock_product'
)
).join(MaterialBase, StockProduct.base_id == MaterialBase.id).filter(
MaterialBase.company_name == company_limit
).subquery()
comp_all = db.session.query(buy_comp.c.outbound_no).union(
db.session.query(semi_comp.c.outbound_no),
db.session.query(prod_comp.c.outbound_no)
).subquery()
stmt = db.session.query(
TransOutbound.outbound_no,
func.max(TransOutbound.outbound_time).label('max_time')
).group_by(TransOutbound.outbound_no)
if keyword_conditions is not None:
stmt = stmt.filter(keyword_conditions)
# ====================================================================
# ★ 高级筛选:父级字段直接过滤,子级字段(SKU/物料名称)走
# 「命中单号子查询 → 按单号 IN」的 EXISTS 语义。
#
# 绝不能写成 stmt.filter(TransOutbound.sku.ilike(...)):
# 那样会在 GROUP BY 前收窄明细范围,展开行里的兄弟明细会丢失。
# ====================================================================
if advanced_filters:
from app.utils.advanced_filter import (
build_predicate, apply_child_condition,
)
# 父级字段:单号/操作人本身就在流水表上,直接 filter(否定操作符
# 走标准 SQL != / NOT LIKE 即可,语义无歧义)
parent_field_map = {
'no': TransOutbound.outbound_no,
'operator': TransOutbound.operator_name,
'consumer_name': TransOutbound.consumer_name,
'outbound_type': TransOutbound.outbound_type,
}
# 子级字段:SKU 直接列 + 物料名称(需三表联查)
child_field_map = {'sku': TransOutbound.sku}
material_stock_models = [
(StockBuy, 'stock_buy'),
(StockSemi, 'stock_semi'),
(StockProduct, 'stock_product'),
]
for cond in advanced_filters:
field = cond.get('field')
if field in parent_field_map:
p = build_predicate(cond, parent_field_map)
if p is not None:
stmt = stmt.filter(p)
continue
if field in child_field_map or field == 'material_name':
# ★ 正/负操作符语义分派:否定 → 整单排除(NOT IN)
stmt = apply_child_condition(
stmt, TransOutbound.outbound_no, TransOutbound, cond,
child_field_map, material_stock_models,
)
continue
if start_date and end_date:
stmt = stmt.filter(TransOutbound.outbound_time.between(start_date, end_date))
# 【行级数据隔离】应用公司过滤到主查询
if company_limit is not None:
stmt = stmt.filter(TransOutbound.outbound_no.in_(comp_all))
# ★ 数据权限:普通用户只看“领用人=本人姓名(不含账号前缀)”的出库记录;
# 同时兼容库里存成“姓名/xiaolongxia”全名的记录(姓名 + '/' 前缀也命中)
if consumer_name:
# 注意:or_ 已在模块顶部导入,此处绝不可再写 `from sqlalchemy import or_`
# —— 函数内出现对 or_ 的赋值(import 即赋值)会让 Python 把 or_ 视为
# 整个函数的局部变量,导致本函数中**位于该行之前**的所有 or_ 调用
# (keyword 搜索分支)抛 UnboundLocalError: referenced before assignment。
_own_out_nos = (
db.session.query(TransOutbound.outbound_no)
.filter(or_(
TransOutbound.consumer_name == consumer_name,
TransOutbound.consumer_name.like(f"{consumer_name}/%")
))
.distinct()
.subquery()
)
stmt = stmt.filter(TransOutbound.outbound_no.in_(_own_out_nos))
stmt = stmt.order_by(desc('max_time'))
# 使用 distinct 确保跨表查询不重复
stmt = stmt.distinct()
pagination = stmt.paginate(page=page, per_page=per_page, error_out=False)
outbound_nos = [row.outbound_no for row in pagination.items]
if not outbound_nos:
return {
'items': [],
'total': 0,
'pages': 0,
'current_page': page
}
# 2. 查询详细记录
details = TransOutbound.query.filter(TransOutbound.outbound_no.in_(outbound_nos)).all()
# 3. 组装数据并查询物品详情
grouped_map = {}
# 映射表模型以便查询
model_map = {
'stock_buy': StockBuy,
'stock_semi': StockSemi,
'stock_product': StockProduct
}
# ==========================================
# ★ 优化步骤 1:第一遍循环,单纯收集所有的 stock_id
# ==========================================
stock_ids_by_table = {'stock_buy': set(), 'stock_semi': set(), 'stock_product': set()}
for d in details:
if d.source_table in stock_ids_by_table and d.stock_id:
stock_ids_by_table[d.source_table].add(d.stock_id)
# ==========================================
# ★ 优化步骤 2:发起批量查询,并强制 JOIN 基础物料表
# ==========================================
# 格式: { ('stock_buy', 101): stock_obj, ... }
preloaded_stocks = {}
for table_name, ids in stock_ids_by_table.items():
if not ids:
continue
ModelClass = model_map[table_name]
# 魔法在这里:in_() 一次性查出所有库存,joinedload 顺便把 base 表的数据一起拉回来
items = ModelClass.query.options(
joinedload(ModelClass.base)
).filter(ModelClass.id.in_(ids)).all()
for item in items:
preloaded_stocks[(table_name, item.id)] = item
# ==========================================
# ★ 优化步骤 3:第二遍循环,纯内存拼装(极速)
# ==========================================
for d in details:
ono = d.outbound_no
if ono not in grouped_map:
grouped_map[ono] = {
'outbound_no': ono,
'outbound_time': d.outbound_time.strftime('%Y-%m-%d %H:%M:%S'),
'outbound_type': d.outbound_type,
'consumer_name': d.consumer_name,
'operator_name': d.operator_name,
'signature_path': d.signature_path,
'remark': d.remark,
'total_amount': 0.0,
'items': []
}
# --- 直接从内存字典中获取,O(1) 复杂度,绝对不触发 SQL ---
item_name, item_spec, item_cat, item_type, batch_sn = "未知物品", "", "", "", "-"
stock_item = preloaded_stocks.get((d.source_table, d.stock_id))
if stock_item:
batch_sn = getattr(stock_item, 'batch_number', None) or getattr(stock_item, 'serial_number', None) or '-'
# 因为前面用了 joinedload,这里调用 .base 瞬间返回,不会去查数据库
if stock_item.base:
item_name = stock_item.base.name
item_spec = stock_item.base.spec_model
item_cat = stock_item.base.category
item_type = stock_item.base.material_type
# 计算金额
price = float(d.unit_price) if d.unit_price else 0
qty = float(d.quantity)
subtotal = price * qty
grouped_map[ono]['total_amount'] += subtotal
grouped_map[ono]['items'].append({
'sku': d.sku,
'name': item_name,
'spec_model': item_spec,
'category': item_cat,
'material_type': item_type,
'quantity': qty,
'unit_price': price,
'subtotal': subtotal,
'batch_sn': batch_sn,
# [新增] 出库时的库位快照(展示用)
'warehouse_location': getattr(d, 'warehouse_location', None) or ''
})
# 4. 排序输出
result_list = []
for ono in outbound_nos:
if ono in grouped_map:
obj = grouped_map[ono]
obj['items'].sort(key=lambda x: x['unit_price'], reverse=True)
obj['total_amount'] = round(obj['total_amount'], 2)
result_list.append(obj)
return {
'items': result_list,
'total': pagination.total,
'pages': pagination.pages,
'current_page': page
}
class OutboundApprovalService:
"""出库审批服务"""
@staticmethod
def generate_request_no():
"""
生成审批单号: APR-OUT-yyyyMMdd-HHmm-当日流水(4位)
"""
beijing_tz = timezone(timedelta(hours=8))
now = datetime.now(beijing_tz)
date_str = now.strftime('%Y%m%d')
time_str = now.strftime('%H%M')
prefix = f"APR-OUT-{date_str}-"
from app.models.outbound import OutboundApproval
latest = db.session.query(OutboundApproval.request_no).filter(
OutboundApproval.request_no.like(f"{prefix}%")
).order_by(OutboundApproval.id.desc()).first()
if latest:
last_seq = int(latest[0].split('-')[-1])
sequence = last_seq + 1
else:
sequence = 1
return f"APR-OUT-{date_str}-{time_str}-{sequence:04d}"
@staticmethod
def _items_require_approval(items):
"""
明细中任一物料命中“出库/借库需审批”→ 需要审批。
申请明细只存 name/spec_model,故按 (name, spec_model) 反查启用物料判定。
"""
from app.models.base import MaterialBase
seen = set()
for item in items:
name = str(item.get('name') or '').strip()
spec = str(item.get('spec_model') or '').strip()
if not name:
continue
key = (name, spec)
if key in seen:
continue
seen.add(key)
hit = MaterialBase.query.filter(
MaterialBase.name == name,
MaterialBase.spec_model == spec,
MaterialBase.is_enabled == True,
MaterialBase.is_approval_required == True
).first()
if hit:
return True
return False
@staticmethod
def create_request(applicant_id, items, allowed_approvers, remark=None, approver_id=None, outbound_type=None, force_approval=False):
"""
创建出库审批单(申请阶段,直接存储前端传来的物料信息快照,不关联具体库存记录)
Args:
applicant_id: 申请人ID
items: 出库物品明细列表,每个物品应包含:
- name: 物料名称 (必填)
- spec_model: 规格型号 (必填)
- quantity: 计划出库数量 (必填)
- warehouse_location: 库位 (可选)
- remark: 物品备注 (可选)
allowed_approvers: 允许审批的人员/角色列表
approver_id: 指定审批人ID(可选,传则覆盖 allowed_approvers)
remark: 申请说明
Returns:
OutboundApproval 实例
Raises:
ValueError: 当 items 为空或缺少必填字段时抛出
"""
from app.models.outbound import OutboundApproval
# 校验 items 非空
if not items:
raise ValueError("出库物品列表不能为空")
# 校验每个物品的宏观字段 (name, spec_model, quantity)
required_fields = ['name', 'spec_model', 'quantity']
for idx, item in enumerate(items):
missing_fields = [f for f in required_fields if f not in item or str(item.get(f) or '').strip() == '']
if missing_fields:
raise ValueError(
f"第 {idx + 1} 条物品缺少必填字段: {', '.join(missing_fields)}。"
f"必须包含: name, spec_model, quantity"
)
try:
qty = float(item.get('quantity', 0))
if qty <= 0:
raise ValueError(f"第 {idx + 1} 条物品的出库数量必须大于0")
except (TypeError, ValueError) as e:
raise ValueError(f"第 {idx + 1} 条物品的 quantity 格式无效: {str(e)}")
# ★ 需审批判定:库管代建(force_approval)或明细含需审批物料 → 走审批;否则默认自动通过
from app.services.approval_control import resolve_approval_control
need_approval, flagged_materials = resolve_approval_control(items)
if force_approval:
need_approval = True
if need_approval and not approver_id:
if force_approval:
raise ValueError("库管代建出库申请必须选择审批人后再提交")
_names = ";".join(f"{m['name']}({m['spec_model'] or '-'})" for m in flagged_materials)
raise ValueError(f"以下物料需审批出库/借库:{_names}。请选择审批人后再提交")
# ★ 指定审批人模式:approver_id 覆盖 allowed_approvers
if approver_id:
allowed_approvers = [{"type": "user", "value": int(approver_id)}]
elif not need_approval:
allowed_approvers = [] # 免审批单不绑定审批人
request_no = OutboundApprovalService.generate_request_no()
if need_approval:
approval = OutboundApproval(
request_no=request_no,
applicant_id=applicant_id,
remark=remark,
outbound_type=outbound_type, # 申请时确定的出库类型,扫码出库时带出
status=0, # 待审批
)
else:
# 默认不审批:创建即已通过(status=1),直接进入“待库管执行”
# ★ approved_at 必须与 created_at(beijing_time) 同为 naive 本地时间:
# approved_at 列是 timestamp without time zone,传 aware 时间会被驱动转成 UTC 存库,导致早 8 小时
approval = OutboundApproval(
request_no=request_no,
applicant_id=applicant_id,
remark=remark,
outbound_type=outbound_type,
status=1,
actual_approver_id=applicant_id,
approved_at=beijing_time(),
)
# ==================================================================
# ★ Phase 1:库存预占
#
# 改造前此处只存「要什么、要多少」的快照,不绑库存行,真正的
# available_quantity 扣减留到执行阶段 —— 于是多张申请可以同时
# claim 同一批货,等到扫码时才发现已被领走(超卖)。
#
# 现在提交即预占:分配器把需求落到具体库存行、立即扣减
# available_quantity,并把 (stock_id, source_table, allocated_qty)
# 连同 reserved 标记写进 items_json,驳回/撤回时据此原样归还。
# ==================================================================
from app.services.inventory_reservation import reserve_for_items
from app.utils.decorators import get_current_company_filter
reserved_items, _shortages = reserve_for_items(
items, company_limit=get_current_company_filter(), strict=True,
)
approval.set_items(reserved_items)
if allowed_approvers:
approval.set_allowed_approvers(allowed_approvers)
else:
approval.allowed_approvers = '[]'
db.session.add(approval)
db.session.commit()
# 仅需审批单通知审批人;免审批单静默进入“待库管执行”
if need_approval:
OutboundApprovalService._notify_new_request(approval, applicant_id, approver_id=approver_id)
return approval
@staticmethod
def _get_emails_by_identifiers(applicant_id=None, role_codes=None):
"""
根据用户ID或角色列表查询邮箱地址
Args:
applicant_id: 用户ID (按 SysUser.id 查找)
role_codes: 角色代码列表,如 ['ADMIN', 'WAREHOUSE_ADMIN']
Returns:
去重后的邮箱地址列表
"""
emails = []
if applicant_id:
user = SysUser.query.get(int(applicant_id))
if user and user.email:
emails.append(user.email)
if role_codes:
for code in role_codes:
users = SysUser.query.filter_by(role=code).all()
for u in users:
if u.email:
emails.append(u.email)
return list(set(emails))
@staticmethod
def _notify_new_request(approval, applicant_id, approver_id=None):
"""发送新申请通知邮件给审批人和申请人(静默处理,不阻断主流程)"""
try:
from flask import current_app
from app.utils.email_service import send_outbound_new_request_notify
from app.models.system import SysUser
applicant_name = ''
applicant_emails = []
# 1. 收集申请人信息
if applicant_id:
user = SysUser.query.get(int(applicant_id))
if user and user.email:
applicant_emails.append(user.email)
applicant_name = str(user.username).split('/')[0] if '/' in (user.username or '') else (user.username or str(applicant_id))
# 2. 收集审批人信息
approver_emails = []
if approver_id:
user = SysUser.query.get(int(approver_id))
if user and user.email:
approver_emails.append(user.email)
else:
# 兜底:按角色查询
approvers = approval.get_allowed_approvers()
role_codes = []
for a in approvers:
if a.get('type') == 'role':
role_codes.append(a.get('value', ''))
approver_emails = OutboundApprovalService._get_emails_by_identifiers(role_codes=role_codes)
# 去重
all_emails = list(set(applicant_emails + approver_emails))
if not all_emails:
current_app.logger.info(f"[Email] 审批单 {approval.request_no} 无收件人邮箱,跳过通知")
return
# 3. 获取物料明细
items = approval.get_items()
# 4. 分别发送邮件
if applicant_emails:
try:
send_outbound_new_request_notify(
to_emails=applicant_emails,
request_no=approval.request_no,
applicant_name=applicant_name,
remark=f"您的出库申请已提交,等待审批。{approval.remark or ''}",
items=items,
is_applicant_notify=True
)
except Exception as e:
current_app.logger.error(f"[Email] 通知申请人失败: {e}")
if approver_emails:
try:
send_outbound_new_request_notify(
to_emails=approver_emails,
request_no=approval.request_no,
applicant_name=applicant_name,
remark=approval.remark or '',
items=items,
is_applicant_notify=False
)
except Exception as e:
current_app.logger.error(f"[Email] 通知审批人失败: {e}")
except Exception as e:
try:
from flask import current_app
current_app.logger.error(f"[Email] 发送新申请通知邮件失败: {e}")
except RuntimeError:
import logging
logging.getLogger(__name__).error(f"[Email] 发送新申请通知邮件失败: {e}")
@staticmethod
def can_approve(approval, user_id, user_role):
"""
检查用户是否有权限审批
Args:
approval: OutboundApproval 实例
user_id: 用户ID
user_role: 用户角色
Returns:
bool, 是否有权限
"""
approvers = approval.get_allowed_approvers()
# 超级管理员可以直接审批
if user_role and user_role.upper() == 'SUPER_ADMIN':
return True
# 拥有 outbound_approval:operation 权限的用户可以审批任意出库单
# (前端按钮以此权限控制显示,后端也必须对齐)
from app.models.system import SysRolePermission
has_op = SysRolePermission.query.filter(
SysRolePermission.role_code == user_role,
SysRolePermission.target_code.in_(['outbound_approval:operation', 'outbound_approval:*'])
).first() is not None
if has_op:
return True
for approver in approvers:
approver_type = approver.get('type', '')
approver_value = approver.get('value', '')
if approver_type == 'user' and str(approver_value) == str(user_id):
return True
if approver_type == 'role' and approver_value == user_role:
return True
return False
@staticmethod
def approve(request_id, user_id, user_role, action='approve', reject_reason=None):
"""
执行审批操作
Args:
request_id: 审批单ID
user_id: 审批人ID
user_role: 审批人角色
action: 'approve' 通过, 'reject' 驳回
reject_reason: 驳回原因
Returns:
(success: bool, message: str, approval: OutboundApproval or None)
"""
from app.models.outbound import OutboundApproval
# ★ 统一取 naive 北京时间,与 created_at(beijing_time) 同口径
current_time = beijing_time()
approval = OutboundApproval.query.get(request_id)
if not approval:
return False, "审批单不存在", None
if approval.status != 0:
return False, f"审批单状态已更新,无法重复审批 (当前状态: {approval.status})", None
if not OutboundApprovalService.can_approve(approval, user_id, user_role):
return False, "您没有审批此单的权限", None
try:
if action == 'approve':
approval.status = 1 # 已通过
approval.actual_approver_id = user_id
approval.approved_at = current_time
# 通过后预占继续保留:货已被本单锁定,直到扫码执行时才释放并扣减
elif action == 'reject':
# ★ Phase 2:驳回即释放预占,把 available_quantity 还回池子,
# 否则这批货会被一张永远不会执行的单永久占住。
from app.services.inventory_reservation import release_reserved
release_reserved(approval.get_items())
approval.status = 2 # 已驳回
approval.reject_reason = reject_reason
else:
return False, "无效的审批操作", None
db.session.commit()
# ★ 审批成功后,发送邮件通知仓库管理员
OutboundApprovalService._notify_approval_result(approval, user_id, action)
return True, "审批成功", approval
except Exception as e:
db.session.rollback()
return False, f"审批失败: {str(e)}", None
@staticmethod
def _notify_approval_result(approval, approver_id, action):
"""发送审批结果通知邮件(静默处理,不阻断主流程)"""
import logging
logger = logging.getLogger(__name__)
try:
from app.utils.email_service import send_outbound_approval_result_notify, send_outbound_dispatch_notify
from app.models.system import SysUser as SU
# 1. 提取申请人信息(供两个分支使用)
applicant_name = ''
applicant_emails = []
if approval.applicant_id:
user = SU.query.get(approval.applicant_id)
if user:
applicant_name = str(user.username).split('/')[0] if '/' in (user.username or '') else (user.username or '')
if user.email:
applicant_emails.append(user.email)
# 2. 提取物料明细(供通过分支使用)
items = approval.get_items() if approval else []
# 3. 分支逻辑
if action == 'approve':
# 3.1 通知库管(带明细)
warehouse_role_codes = ['WAREHOUSE_MGR', 'OUTBOUND']
warehouse_emails = OutboundApprovalService._get_emails_by_identifiers(role_codes=warehouse_role_codes)
if warehouse_emails:
try:
send_outbound_dispatch_notify(
to_emails=warehouse_emails,
request_no=approval.request_no,
applicant_name=applicant_name,
items=items
)
except Exception as e:
logger.error(f"[Email] 通知库管失败: {e}")
# 3.2 通知申请人(审批通过,带完整物料清单)
if applicant_emails:
try:
send_outbound_dispatch_notify(
to_emails=applicant_emails,
request_no=approval.request_no,
applicant_name=applicant_name,
items=items
)
except Exception as e:
logger.error(f"[Email] 通知申请人(通过)失败: {e}")
elif action == 'reject':
# 3.3 通知申请人(已驳回)
if applicant_emails:
try:
send_outbound_approval_result_notify(
to_emails=applicant_emails,
request_no=approval.request_no,
is_passed=False,
reject_reason=approval.reject_reason or '未说明原因',
applicant_name=applicant_name
)
except Exception as e:
logger.error(f"[Email] 通知申请人驳回失败: {e}")
else:
logger.warning("[Email] 申请人无邮箱,无法发送驳回通知")
except Exception as e:
import traceback
traceback.print_exc()
logger.error(f"[Email] 外层发送异常: {e}")
@staticmethod
def close_request(request_id, user_id, user_role):
"""
手动完结/作废审批单(状态 1-已通过 → 4-已完结)
适用场景:已通过但无法出库/作废的单据,库管手动清理,
使其从"已审批通过"列表中消失。
Args:
request_id: 审批单ID
user_id: 操作人ID
user_role: 操作人角色
Returns:
(success: bool, message: str, approval: OutboundApproval or None)
"""
from app.models.outbound import OutboundApproval
approval = OutboundApproval.query.get(request_id)
if not approval:
return False, "审批单不存在", None
if approval.status != 1:
return False, f"仅「已通过」的审批单可完结 (当前状态: {approval.status})", None
# 权限检查:超级管理员、审批人 或 拥有出库操作权限(库管/主管)
if not OutboundApprovalService.can_approve(approval, user_id, user_role):
# 放宽:拥有 outbound_create:operation 的用户(库管)也可完结
from app.models.system import SysRolePermission
has_outbound_op = SysRolePermission.query.filter(
SysRolePermission.role_code == user_role,
SysRolePermission.target_code.in_(['outbound_create:operation', 'outbound_create:*'])
).first() is not None
if not has_outbound_op:
return False, "您没有完结此单的权限", None
try:
approval.status = 4 # 4-已完结(手动作废)
approval.actual_approver_id = user_id
approval.approved_at = None
db.session.commit()
return True, "审批单已完结", approval
except Exception as e:
db.session.rollback()
return False, f"完结失败: {str(e)}", None
def get_request_list(page=1, per_page=10, applicant_id=None, status=None):
"""
获取审批单列表
Args:
page: 页码
per_page: 每页数量
applicant_id: 按申请人筛选 (可选)
status: 按状态筛选 (可选)
Returns:
分页结果
"""
from app.models.outbound import OutboundApproval
from app.models.system import SysUser
from app.utils.decorators import get_current_company_filter
from sqlalchemy import desc
query = OutboundApproval.query
if applicant_id:
query = query.filter(OutboundApproval.applicant_id == applicant_id)
if status is not None:
query = query.filter(OutboundApproval.status == status)
# 【行级数据隔离】通过申请人关联用户表过滤公司
company_limit = get_current_company_filter()
if company_limit is not None:
query = query.join(SysUser, OutboundApproval.applicant_id == SysUser.id) \
.filter(SysUser.department == company_limit)
query = query.order_by(desc(OutboundApproval.created_at))
pagination = query.paginate(page=page, per_page=per_page, error_out=False)
return {
'items': [item.to_dict() for item in pagination.items],
'total': pagination.total,
'pages': pagination.pages,
'current_page': page
}
@staticmethod
def get_request_by_id(request_id):
"""根据ID获取审批单"""
from app.models.outbound import OutboundApproval
return OutboundApproval.query.get(request_id)