Files
goods/tools/bypos-collector/main.go
ceyhandagdas51272 b08495b6f6
CI / Go (api) (pull_request) Successful in 55s
CI / Python (ingestion) (pull_request) Failing after 31s
CI / Migrations (postgres) (pull_request) Failing after 32s
feat: 持久化采集记录去重 + Goods系统批量导入
- collected.txt 跨文件/跨次记录已查询条码,避免重复采集
- Web UI 显示历史已采集条码计数
- Goods API 新增 POST /api/import/bypos 批量导入端点(含upsert)
- 采集器 Web UI 新增一键导入按钮(登录+上传JSONL)
- Windows exe 重新编译

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-24 04:40:26 +00:00

274 lines
7.5 KiB
Go

package main
import (
"bufio"
"bytes"
"embed"
"encoding/csv"
"encoding/json"
"flag"
"fmt"
"io"
"io/fs"
"log"
"net"
"net/http"
"os"
"os/exec"
"runtime"
"strings"
"time"
)
//go:embed web/*
var webFS embed.FS
var collector *Collector
// waitExit keeps the console window open on Windows so the user can read errors.
func waitExit() {
if runtime.GOOS == "windows" {
fmt.Println("\n按回车键退出...")
bufio.NewReader(os.Stdin).ReadBytes('\n')
}
}
func main() {
// Log to file so crashes are diagnosable even if console closes.
lf, lfErr := os.OpenFile("bypos-collector.log", os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644)
if lfErr == nil {
log.SetOutput(io.MultiWriter(os.Stderr, lf))
defer lf.Close()
}
defer func() {
if r := recover(); r != nil {
log.Printf("程序崩溃: %v", r)
waitExit()
}
}()
log.Println("bypos-collector 启动中...")
addr := flag.String("addr", "127.0.0.1:8765", "本地监听地址")
noOpen := flag.Bool("no-open", false, "不自动打开浏览器")
sdog := flag.String("sdogid", "", "中心库账号 id(优先读 $BYPOS_SDOGID 环境变量)")
flag.Parse()
collector = NewCollector(*sdog)
if collector.sdogID == "" {
log.Println("警告: 未配置 sdogid,请通过 $BYPOS_SDOGID 环境变量、-sdogid 参数或控制台输入框指定")
}
sub, err := fs.Sub(webFS, "web")
if err != nil {
log.Printf("错误: 无法加载内嵌 web 资源: %v", err)
waitExit()
return
}
mux := http.NewServeMux()
mux.Handle("/", http.FileServer(http.FS(sub)))
mux.HandleFunc("/api/start", handleStart)
mux.HandleFunc("/api/stop", handleStop)
mux.HandleFunc("/api/stats", handleStats)
mux.HandleFunc("/api/download", handleDownload)
mux.HandleFunc("/api/export.csv", handleExportCSV)
mux.HandleFunc("/api/import", handleImport)
ln, err := net.Listen("tcp", *addr)
if err != nil {
log.Printf("端口 %s 被占用,自动选择可用端口...", *addr)
ln, err = net.Listen("tcp", "127.0.0.1:0")
if err != nil {
log.Printf("错误: 无法监听: %v", err)
waitExit()
return
}
}
realAddr := ln.Addr().String()
urlStr := "http://" + realAddr + "/"
fmt.Println("==============================================")
fmt.Println(" 中心库商品采集器 bypos-collector")
fmt.Println(" 控制台: " + urlStr)
fmt.Println(" 关闭本窗口即停止程序")
fmt.Println("==============================================")
log.Printf("监听地址: %s", realAddr)
if !*noOpen {
go openBrowser(urlStr)
}
if err := http.Serve(ln, mux); err != nil {
log.Printf("HTTP 服务异常退出: %v", err)
waitExit()
}
}
func openBrowser(url string) {
time.Sleep(600 * time.Millisecond)
var cmd string
var args []string
switch runtime.GOOS {
case "windows":
cmd = "rundll32"
args = []string{"url.dll,FileProtocolHandler", url}
case "darwin":
cmd = "open"
args = []string{url}
default:
cmd = "xdg-open"
args = []string{url}
}
_ = exec.Command(cmd, args...).Start()
}
func writeJSON(w http.ResponseWriter, code int, v interface{}) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(code)
json.NewEncoder(w).Encode(v)
}
func handleStart(w http.ResponseWriter, r *http.Request) {
if r.Method != "POST" {
http.Error(w, "method", 405)
return
}
var req JobReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, 400, map[string]string{"error": "请求格式错误"})
return
}
if err := collector.Start(req); err != nil {
writeJSON(w, 400, map[string]string{"error": err.Error()})
return
}
writeJSON(w, 200, map[string]string{"ok": "started"})
}
func handleStop(w http.ResponseWriter, r *http.Request) {
collector.Stop()
writeJSON(w, 200, map[string]string{"ok": "stopping"})
}
func handleStats(w http.ResponseWriter, r *http.Request) {
writeJSON(w, 200, map[string]interface{}{
"stats": collector.snapshot(),
"recent": collector.recentResults(),
})
}
func handleDownload(w http.ResponseWriter, r *http.Request) {
s := collector.snapshot()
if s.OutFile == "" {
http.Error(w, "no output yet", 404)
return
}
f, err := os.Open(s.OutFile)
if err != nil {
http.Error(w, err.Error(), 404)
return
}
defer f.Close()
w.Header().Set("Content-Type", "application/x-ndjson; charset=utf-8")
w.Header().Set("Content-Disposition", "attachment; filename=products.jsonl")
io.Copy(w, f)
}
func handleImport(w http.ResponseWriter, r *http.Request) {
if r.Method != "POST" {
http.Error(w, "method", 405)
return
}
var req struct {
APIURL string `json:"api_url"`
Username string `json:"username"`
Password string `json:"password"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, 400, map[string]string{"error": "请求格式错误"})
return
}
req.APIURL = strings.TrimRight(req.APIURL, "/")
if req.APIURL == "" || req.Username == "" || req.Password == "" {
writeJSON(w, 400, map[string]string{"error": "请填写完整的 API 地址、用户名和密码"})
return
}
s := collector.snapshot()
if s.OutFile == "" {
writeJSON(w, 400, map[string]string{"error": "没有采集数据可导入"})
return
}
client := &http.Client{Timeout: 60 * time.Second}
loginBody, _ := json.Marshal(map[string]string{"username": req.Username, "password": req.Password})
loginResp, err := client.Post(req.APIURL+"/api/login", "application/json", bytes.NewReader(loginBody))
if err != nil {
writeJSON(w, 502, map[string]string{"error": "无法连接 Goods 系统: " + err.Error()})
return
}
defer loginResp.Body.Close()
var loginResult struct {
Token string `json:"token"`
Error *struct {
Message string `json:"message"`
} `json:"error"`
}
json.NewDecoder(loginResp.Body).Decode(&loginResult)
if loginResp.StatusCode != 200 || loginResult.Token == "" {
msg := "登录失败"
if loginResult.Error != nil {
msg = loginResult.Error.Message
}
writeJSON(w, 401, map[string]string{"error": msg})
return
}
jsonlData, err := os.ReadFile(s.OutFile)
if err != nil {
writeJSON(w, 500, map[string]string{"error": "读取采集文件失败: " + err.Error()})
return
}
importReq, _ := http.NewRequest("POST", req.APIURL+"/api/import/bypos", bytes.NewReader(jsonlData))
importReq.Header.Set("Content-Type", "application/x-ndjson")
importReq.Header.Set("Authorization", "Bearer "+loginResult.Token)
importResp, err := client.Do(importReq)
if err != nil {
writeJSON(w, 502, map[string]string{"error": "导入请求失败: " + err.Error()})
return
}
defer importResp.Body.Close()
respBody, _ := io.ReadAll(importResp.Body)
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(importResp.StatusCode)
w.Write(respBody)
}
func handleExportCSV(w http.ResponseWriter, r *http.Request) {
s := collector.snapshot()
if s.OutFile == "" {
http.Error(w, "no output yet", 404)
return
}
f, err := os.Open(s.OutFile)
if err != nil {
http.Error(w, err.Error(), 404)
return
}
defer f.Close()
w.Header().Set("Content-Type", "text/csv; charset=utf-8")
w.Header().Set("Content-Disposition", "attachment; filename=products.csv")
w.Write([]byte{0xEF, 0xBB, 0xBF}) // UTF-8 BOM so Excel reads Chinese correctly
cw := csv.NewWriter(w)
cw.Write([]string{"barcode", "name", "spec", "unit", "area", "manufacturer", "license", "in_price", "sell_price", "status", "retmsg", "fetched_at", "source"})
dec := json.NewDecoder(f)
for {
var p Product
if err := dec.Decode(&p); err != nil {
break
}
cw.Write([]string{p.Barcode, p.Name, p.Spec, p.Unit, p.Area, p.Manufacturer, p.License, p.InPrice, p.SellPrice, p.Status, p.RetMsg, p.FetchedAt, p.Source})
}
cw.Flush()
}