Files
bj_power/bj_power_mes/internal/db/db.go
T
SunYF 05b0b8a909 feat: 完成客户第12条需求+多系统改造
1. 新增预警表添加工单号字段,支持按工单号检索预警
2. 添加工单过滤在制品/绩效报表/预警中心查询
3. 补全WMS台账字段口径与权限配置
4. 修复备料单与数量不符处理流程
5. 新增配送进度/库位联动/自动出库功能
6. 修复AGV配送状态更新逻辑
7. 优化退料与库存管理细节
2026-09-19 17:04:58 +08:00

243 lines
12 KiB
Go
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.
package db
import (
"context"
"database/sql"
"fmt"
"time"
"bj_power_mes/ent"
"entgo.io/ent/dialect"
entsql "entgo.io/ent/dialect/sql"
_ "github.com/jackc/pgx/v5/stdlib"
)
// DatabaseConf 数据库配置,DSN 由调用方从 yaml 读取
type DatabaseConf struct {
Host string `json:",default=127.0.0.1"`
Port int `json:",default=5432"`
User string `json:",default=postgres"`
Password string `json:",default=postgres"`
Dbname string `json:",default=bj_power_mes"`
MaxIdle int `json:",default=10"`
MaxOpen int `json:",default=20"`
}
// DSN 从配置拼装连接串(不在代码中硬编码口令)
func (c DatabaseConf) DSN() string {
return fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=%s sslmode=disable",
c.Host, c.Port, c.User, c.Password, c.Dbname)
}
// MustNewDB 打开 PostgreSQL 连接并包装为 ent.Client
func MustNewDB(c DatabaseConf) (*ent.Client, *sql.DB) {
sqlDB, err := sql.Open("pgx", c.DSN())
if err != nil {
panic(fmt.Sprintf("open db failed: %v", err))
}
sqlDB.SetMaxIdleConns(c.MaxIdle)
sqlDB.SetMaxOpenConns(c.MaxOpen)
sqlDB.SetConnMaxLifetime(time.Hour)
drv := entsql.OpenDB(dialect.Postgres, sqlDB)
client := ent.NewClient(ent.Driver(drv))
return client, sqlDB
}
// NewDB 打开连接,失败返回 error(用于 migrate 等工具)
func NewDB(c DatabaseConf) (*sql.DB, error) {
sqlDB, err := sql.Open("pgx", c.DSN())
if err != nil {
return nil, err
}
if err := sqlDB.Ping(); err != nil {
return nil, err
}
return sqlDB, nil
}
// EnsureDB 若目标数据库不存在则自动创建。
// 先连接 PostgreSQL 内置维护库 maintenance db(默认 postgres),检查目标库,
// 不存在则执行 CREATE DATABASE,便于重置/建表命令在空环境一键运行。
func EnsureDB(c DatabaseConf) error {
maintenanceDSN := fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=postgres sslmode=disable",
c.Host, c.Port, c.User, c.Password)
sqlDB, err := sql.Open("pgx", maintenanceDSN)
if err != nil {
return err
}
defer sqlDB.Close()
var exists bool
if err := sqlDB.QueryRow(
`SELECT EXISTS(SELECT 1 FROM pg_database WHERE datname = $1)`, c.Dbname,
).Scan(&exists); err != nil {
return err
}
if exists {
return nil
}
// 库名需要用双引号包裹以支持大写/特殊字符
if _, err := sqlDB.Exec(`CREATE DATABASE "` + c.Dbname + `"`); err != nil {
return err
}
return nil
}
// schemaPatchSQL 幂等列补丁:ent Schema.Create 只建缺失的新表,绝不会给既有表加列;
// 因此对"已有表新增列"的场景必须走显式 ALTER。这里统一维护补丁清单,migrate/seed/reset-all 三入口共用。
var schemaPatchSQL = []string{
// BOM 料增加"装配工序"维度(0=不参与工位绑定校验;1~12=在第几道工序装配)
// 注意:BomItem schema 表名注解为 work_order_bom(非默认名)
`ALTER TABLE work_order_bom ADD COLUMN IF NOT EXISTS process_code integer NOT NULL DEFAULT 0`,
// 第三批:工程编号三级归属(合同号 → 工程编号 → 产品序号)
`ALTER TABLE work_order ADD COLUMN IF NOT EXISTS contract_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE work_order ADD COLUMN IF NOT EXISTS project_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE work_order ADD COLUMN IF NOT EXISTS product_serial varchar(64) NOT NULL DEFAULT ''`,
// P0-4:过程巡检按步骤提交(步骤维度字段)
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS process_code integer NOT NULL DEFAULT 0`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS step_id integer NOT NULL DEFAULT 0`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS step_name varchar(128) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS measured_value varchar(64) NOT NULL DEFAULT ''`,
// 上生产线前剩余项(H 质量检验 / I 产品物料清单 / J 步骤附件 / M 工单字段·绩效)
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS report_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS material_code varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS material_name varchar(128) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS manufacturer varchar(128) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS inspection_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS attachment_ids varchar(512) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS disposal_type varchar(20) NOT NULL DEFAULT ''`,
`ALTER TABLE inspection_record ADD COLUMN IF NOT EXISTS disposal_remark varchar(500) NOT NULL DEFAULT ''`,
// work_order_bom 表名注解为 work_order_bom
`ALTER TABLE work_order_bom ADD COLUMN IF NOT EXISTS related_standard varchar(255) NOT NULL DEFAULT ''`,
`ALTER TABLE process_step ADD COLUMN IF NOT EXISTS attachment jsonb NOT NULL DEFAULT '[]'::jsonb`,
`ALTER TABLE work_order ADD COLUMN IF NOT EXISTS created_by varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE work_order ADD COLUMN IF NOT EXISTS due_date timestamp`,
`ALTER TABLE workpiece_process ADD COLUMN IF NOT EXISTS duration_sec integer NOT NULL DEFAULT 0`,
// 客户第12条「记录查看加一项看工单号查询」:预警表补工单号,供预警中心按工单号检索
`ALTER TABLE alert ADD COLUMN IF NOT EXISTS order_no varchar(64) NOT NULL DEFAULT ''`,
// 出库与配送闭环改造(2026-09-19):拆开被污染的账本字段 + 数量不符处理侧 + 补料来源
// material_requestreceived_* 拆出「已接料量」(sent_qty 锁定为已发出量=幂等基准,唯一写者=自动出库)
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS received_qty double precision NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS received_at bigint NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS receive_by varchar(64) NOT NULL DEFAULT ''`,
// 来源标记:PLAN=按日排产备料 / REFILL=工位叫料·数量不符补发(补料优先出库)
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS source varchar(16) NOT NULL DEFAULT 'PLAN'`,
// 超 BOM 需求放行留痕:超量与原因(损耗/报废/返修/其他),自动出库不填=不超领
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS over_qty double precision NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS over_reason varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS ref_report_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS need_agv boolean NOT NULL DEFAULT false`,
// material_qty_report:处理侧字段(补发=生成 REFILL 备料单;退库=WMS 退库单;调整=仅关闭)
`ALTER TABLE material_qty_report ADD COLUMN IF NOT EXISTS handle_type varchar(16) NOT NULL DEFAULT ''`,
`ALTER TABLE material_qty_report ADD COLUMN IF NOT EXISTS handle_doc_no varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE material_qty_report ADD COLUMN IF NOT EXISTS handle_note varchar(255) NOT NULL DEFAULT ''`,
`ALTER TABLE material_qty_report ADD COLUMN IF NOT EXISTS handled_by varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE material_qty_report ADD COLUMN IF NOT EXISTS handled_at bigint NOT NULL DEFAULT 0`,
// 工位内置标识 + 是否有接驳台(实体位置由内置数据给定,页面不可手改)
`ALTER TABLE station ADD COLUMN IF NOT EXISTS is_builtin boolean NOT NULL DEFAULT false`,
`ALTER TABLE station ADD COLUMN IF NOT EXISTS has_dock boolean NOT NULL DEFAULT false`,
// 存量数据回填:种子初始化的工位号(0/1..13)一律为内置工位;
// 其中 1~10 号工位对应实体接驳台(R1/R2),0/11/12/13 无接驳台。
// 页面新增的工位号从 14 起,不受本回填影响。
`UPDATE station SET is_builtin = true WHERE station_no BETWEEN 0 AND 13`,
`UPDATE station SET has_dock = true WHERE station_no BETWEEN 1 AND 10`,
// ---- 本次重构:删除工位组合/接驳台冗余字段,改用关联表 station_process + station.dock_code 驱动 ----
`ALTER TABLE work_order DROP COLUMN IF EXISTS process_seq`,
`ALTER TABLE workpiece DROP COLUMN IF EXISTS process_seq`,
`ALTER TABLE daily_plan DROP COLUMN IF EXISTS dock_codes`,
// station 引用接驳台(与 WMS dock.dock_code 对应,1:1
`ALTER TABLE station ADD COLUMN IF NOT EXISTS dock_code varchar(20) NOT NULL DEFAULT ''`,
// 工位↔接驳台 1:1 唯一约束:一个接驳台只能绑定一个工位、一个工位只能有一个接驳台。
// 部分唯一索引只约束已绑定(dock_code 非空)的行,未绑定工位不参与唯一性。
// 存量若已存在重复绑定会导致建索引失败(仅告警不阻断启动),须先在本页解绑重复项。
`CREATE UNIQUE INDEX IF NOT EXISTS ux_station_dock_code ON station (dock_code) WHERE dock_code <> ''`,
// 备料单工位级 + 分批出库
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS station_no integer NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS process_code integer NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS sent_qty double precision NOT NULL DEFAULT 0`,
`ALTER TABLE material_request ADD COLUMN IF NOT EXISTS batch_no varchar(64) NOT NULL DEFAULT ''`,
// 日排产各工位产量分配
`ALTER TABLE daily_plan ADD COLUMN IF NOT EXISTS station_plan_qty jsonb NOT NULL DEFAULT '{}'::jsonb`,
// ---- 物料档案「简称」删除(2026-09-19 用户裁决)----
// 本地 product_type 仅作 WMS 不可达时的离线兜底缓存:原 remark(备注) 一并删除,
// 改存 spec/unit,与 WMS 物料档案必填口径(图号/名称/规格/单位/品类/类型)保持一致。
`ALTER TABLE product_type ADD COLUMN IF NOT EXISTS spec varchar(128) NOT NULL DEFAULT ''`,
`ALTER TABLE product_type ADD COLUMN IF NOT EXISTS unit varchar(32) NOT NULL DEFAULT ''`,
`ALTER TABLE product_type DROP COLUMN IF EXISTS remark`,
// ---- 附件统一存储方案(2026-09-19):年/月/日/文件类型/uuid.ext ----
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS file_type varchar(32) NOT NULL DEFAULT 'OTHER'`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS file_ext varchar(16) NOT NULL DEFAULT ''`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS file_md5 varchar(64) NOT NULL DEFAULT ''`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS mime_type varchar(128) NOT NULL DEFAULT ''`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS archived boolean NOT NULL DEFAULT false`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS archived_at bigint NOT NULL DEFAULT 0`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS deleted boolean NOT NULL DEFAULT false`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS deleted_at bigint NOT NULL DEFAULT 0`,
`ALTER TABLE attachment ADD COLUMN IF NOT EXISTS deleted_by varchar(64) NOT NULL DEFAULT ''`,
}
// EnsureSchema ent 建表 + 幂等列补丁(只增,不删数据)
func EnsureSchema(client *ent.Client, sqlDB *sql.DB) error {
if err := client.Schema.Create(context.Background()); err != nil {
return err
}
for _, q := range schemaPatchSQL {
if _, err := sqlDB.Exec(q); err != nil {
return fmt.Errorf("schema patch failed: %v", err)
}
}
return nil
}
// AutoMigrate 执行 ent 自动建表/补齐(只增,不删数据)。供命令行 migrate 使用。
func AutoMigrate(c DatabaseConf) error {
client, sqlDB := MustNewDB(c)
defer sqlDB.Close()
return EnsureSchema(client, sqlDB)
}
// DropAllTables 清空 public schema 下所有业务表(破坏性,仅供命令行 reset-all 使用)。
func DropAllTables(c DatabaseConf) error {
sqlDB, err := NewDB(c)
if err != nil {
return err
}
defer sqlDB.Close()
rows, err := sqlDB.Query(`SELECT tablename FROM pg_tables WHERE schemaname='public'`)
if err != nil {
return err
}
var names []string
for rows.Next() {
var t string
if err := rows.Scan(&t); err != nil {
rows.Close()
return err
}
names = append(names, t)
}
rows.Close()
for _, t := range names {
if _, err := sqlDB.Exec(`DROP TABLE IF EXISTS "` + t + `" CASCADE`); err != nil {
return err
}
}
return nil
}