b08495b6f6
- 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>
274 lines
7.5 KiB
Go
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()
|
|
}
|