Files
bj_power/bj_power_mes/internal/db/db.go
T
SunYF 62e45b7379 feat: 完成附件存储重构与磁盘监控功能,同步物料档案与入库单规则更新
1.  新增跨平台磁盘空间监控能力,每日7点自动检测附件目录剩余空间,触发阈值告警
2.  重构附件存储方案为年/月/日/文件类型分层结构,统一MES与WMS的附件管理逻辑
3.  对齐物料档案与入库单的质量状态校验规则,仅合格品计入库存与出库
4.  实现入库单作废功能与区域库位的多级父子结构管理
5.  删除物料简称字段,补充规格型号与单位必填项,统一系统数据口径
6.  新增附件中心与磁盘状态查询接口,完善权限控制与操作日志
2026-09-19 10:54:15 +08:00

221 lines
10 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`,
// 工位内置标识 + 是否有接驳台(实体位置由内置数据给定,页面不可手改)
`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
}