clean: remove deprecated utils/ package
- Remove all 5 files from utils/ (serial.go, socket.go, pulsar.go, sqlite.go, telegram.go) - cmd/root.go: remove utils import and all utils. references - cmd/test.go: rewrite to use serial, storage, transport packages directly with proper error logging via zap
This commit is contained in:
+1
-14
@@ -26,7 +26,6 @@ import (
|
||||
"go.uber.org/zap/zapcore"
|
||||
|
||||
"it2000.com.cn/tele-recv/config"
|
||||
"it2000.com.cn/tele-recv/utils"
|
||||
)
|
||||
|
||||
const ModeName = "TELEGRAM_MODE"
|
||||
@@ -77,7 +76,6 @@ func init() {
|
||||
rootCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
|
||||
|
||||
initLog()
|
||||
utils.Log = logger
|
||||
}
|
||||
|
||||
// initConfig reads in config file and ENV variables if set.
|
||||
@@ -89,7 +87,7 @@ func initConfig() {
|
||||
}
|
||||
appCfg = cfg
|
||||
|
||||
// Populate legacy globals for backward compatibility
|
||||
// Populate legacy globals
|
||||
device = cfg.Serial.Device
|
||||
baudrate = cfg.Serial.Baudrate
|
||||
lograw = cfg.Serial.LogRaw
|
||||
@@ -101,8 +99,6 @@ func initConfig() {
|
||||
name = cfg.Pulsar.Name
|
||||
tcp = cfg.Telegram.TCP
|
||||
pulsar = cfg.Telegram.Pulsar
|
||||
utils.Tcp = tcp
|
||||
utils.Pulsar = pulsar
|
||||
}
|
||||
|
||||
func initLog() {
|
||||
@@ -115,15 +111,6 @@ func initLog() {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
rawFile := "./logs/raw-%Y-%m-%d-%H.txt"
|
||||
utils.RawLog, err = rotatelogs.New(
|
||||
rawFile,
|
||||
rotatelogs.WithMaxAge(60*24*time.Hour),
|
||||
rotatelogs.WithRotationTime(time.Hour))
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
filePriority := zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
|
||||
return lvl >= zapcore.DebugLevel
|
||||
})
|
||||
|
||||
+28
-7
@@ -17,7 +17,10 @@ package cmd
|
||||
|
||||
import (
|
||||
"github.com/spf13/cobra"
|
||||
"it2000.com.cn/tele-recv/utils"
|
||||
"go.uber.org/zap"
|
||||
"it2000.com.cn/tele-recv/serial"
|
||||
"it2000.com.cn/tele-recv/storage"
|
||||
"it2000.com.cn/tele-recv/transport"
|
||||
)
|
||||
|
||||
// testCmd represents the test command
|
||||
@@ -34,17 +37,35 @@ Test load sqlite database and init table`,
|
||||
|
||||
func test() {
|
||||
logger.Info("Testing environment")
|
||||
_ = utils.OpenPort(device, baudrate)
|
||||
|
||||
err := utils.InitDb(dbFile, dbInit)
|
||||
// Test serial port
|
||||
port, err := serial.Open(device, baudrate)
|
||||
if err != nil {
|
||||
println("create database error ", err)
|
||||
logger.Warn("serial port open failed (expected if no device)", zap.Error(err))
|
||||
} else {
|
||||
port.Close()
|
||||
logger.Info("serial port ok")
|
||||
}
|
||||
|
||||
utils.CreateProducer(pulsarUrl, topic, name)
|
||||
utils.ClosePulsar()
|
||||
// Test SQLite
|
||||
store, err := storage.New(dbFile, dbInit)
|
||||
if err != nil {
|
||||
logger.Error("database init failed", zap.Error(err))
|
||||
} else {
|
||||
logger.Info("database ok")
|
||||
store.Close()
|
||||
}
|
||||
|
||||
// Test Pulsar
|
||||
sender, err := transport.NewPulsarSender(pulsarUrl, topic, name)
|
||||
if err != nil {
|
||||
logger.Warn("pulsar connection failed (expected if no broker)", zap.Error(err))
|
||||
} else {
|
||||
logger.Info("pulsar ok")
|
||||
sender.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(testCmd)
|
||||
}
|
||||
}
|
||||
@@ -1,75 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/apache/pulsar-client-go/pulsar"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
var (
|
||||
pulsarClient pulsar.Client
|
||||
producer pulsar.Producer
|
||||
pulsarUrl string
|
||||
pulsarTopic string
|
||||
clientName string
|
||||
)
|
||||
|
||||
func CreateProducer(url string, topic string, name string) {
|
||||
pulsarUrl = url
|
||||
pulsarTopic = topic
|
||||
clientName = name
|
||||
create()
|
||||
}
|
||||
|
||||
func create() {
|
||||
var err error
|
||||
pulsarClient, err = pulsar.NewClient(pulsar.ClientOptions{
|
||||
URL: pulsarUrl,
|
||||
})
|
||||
if err != nil {
|
||||
Log.Error("error in create pulsar client ", zap.Error(err))
|
||||
}
|
||||
producer, err = pulsarClient.CreateProducer(pulsar.ProducerOptions{
|
||||
Topic: pulsarTopic,
|
||||
Name: clientName,
|
||||
})
|
||||
if err != nil {
|
||||
Log.Error("error in create pulsar producer ", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
func resetPulsarProducer() {
|
||||
if producer != nil {
|
||||
producer.Close()
|
||||
}
|
||||
if pulsarClient != nil {
|
||||
pulsarClient.Close()
|
||||
}
|
||||
create()
|
||||
}
|
||||
|
||||
func PulsarSend(data string) error {
|
||||
var err error
|
||||
var id pulsar.MessageID
|
||||
for i := 0; i < 5; i++ {
|
||||
id, err = producer.Send(context.Background(), &pulsar.ProducerMessage{
|
||||
Payload: []byte(data),
|
||||
})
|
||||
if err == nil {
|
||||
Log.Info("send ", zap.Any("id", id))
|
||||
break
|
||||
}
|
||||
resetPulsarProducer()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func ClosePulsar() {
|
||||
if producer != nil {
|
||||
producer.Close()
|
||||
}
|
||||
if pulsarClient != nil {
|
||||
pulsarClient.Close()
|
||||
}
|
||||
}
|
||||
@@ -1,50 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/argandas/serial"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
var (
|
||||
device string
|
||||
baudrate int
|
||||
// port *serial.Port
|
||||
sp *serial.SerialPort
|
||||
portOpened = false
|
||||
Log *zap.Logger
|
||||
)
|
||||
|
||||
const readDuration = 500 * time.Millisecond
|
||||
|
||||
// OpenPort open serial port
|
||||
func OpenPort(device string, baudrate int) error {
|
||||
Log.Info("opening serial device ",
|
||||
zap.String("port", device),
|
||||
zap.Int("baudrate", baudrate))
|
||||
sp = serial.New()
|
||||
sp.EOL('\r')
|
||||
sp.Verbose = false
|
||||
err := sp.Open(device, baudrate, time.Second*3)
|
||||
Log.Info("open ", zap.String("device", device), zap.Error(err))
|
||||
if err == nil {
|
||||
portOpened = true
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
//ReadPort read byte array from port
|
||||
func ReadPort() (string, error) {
|
||||
// buffer := make([]byte, 1)
|
||||
// if !ServerRunning {
|
||||
// return buffer, 0
|
||||
// }
|
||||
|
||||
return sp.ReadLine()
|
||||
}
|
||||
|
||||
//IsPortOpen check serial port status
|
||||
func IsPortOpen() bool {
|
||||
return portOpened
|
||||
}
|
||||
-103
@@ -1,103 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"io"
|
||||
"net"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
const ServerType = "tcp"
|
||||
|
||||
var ServerRunning = false
|
||||
var client net.Conn
|
||||
var server net.Listener
|
||||
var ClientReady = false
|
||||
|
||||
func Listen(address string) {
|
||||
var err error
|
||||
server, err = net.Listen(ServerType, address)
|
||||
if err != nil {
|
||||
Log.Fatal("error in create server ", zap.Error(err))
|
||||
}
|
||||
defer server.Close()
|
||||
Log.Info("listen on ", zap.String("address", address))
|
||||
|
||||
acceptFailCount := 0
|
||||
const maxAcceptFail = 3
|
||||
|
||||
for ServerRunning {
|
||||
conn, err := server.Accept()
|
||||
if err != nil {
|
||||
acceptFailCount++
|
||||
Log.Error("error in create connection ",
|
||||
zap.Error(err),
|
||||
zap.Int("failCount", acceptFailCount))
|
||||
if acceptFailCount >= maxAcceptFail {
|
||||
ServerRunning = false
|
||||
Log.Fatal("accept failed too many times, shutting down")
|
||||
}
|
||||
continue
|
||||
}
|
||||
acceptFailCount = 0
|
||||
|
||||
if client == nil {
|
||||
client = conn
|
||||
ClientReady = false
|
||||
Log.Info("client connected on ", zap.Any("address", conn.RemoteAddr()))
|
||||
LoadUnprocessed()
|
||||
} else {
|
||||
Log.Info("remove old client")
|
||||
closeClient()
|
||||
client = conn
|
||||
}
|
||||
if IsClientConnected() {
|
||||
connCheck()
|
||||
}
|
||||
time.Sleep(2 * time.Second)
|
||||
}
|
||||
}
|
||||
|
||||
func WriteToClient(data string) error {
|
||||
if client != nil {
|
||||
_, err := client.Write([]byte(data + "\r\n"))
|
||||
if err != nil {
|
||||
Log.Error("error in write data to client socket ", zap.Error(err))
|
||||
closeClient()
|
||||
}
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func connCheck() bool {
|
||||
_, err := client.Read(make([]byte, 0))
|
||||
if err != nil && err != io.EOF {
|
||||
// this connection is invalid
|
||||
Log.Error("conn closed....", zap.Error(err))
|
||||
client = nil
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func IsClientConnected() bool {
|
||||
return client != nil
|
||||
}
|
||||
|
||||
func closeClient() {
|
||||
if client != nil {
|
||||
client.Close()
|
||||
client = nil
|
||||
}
|
||||
}
|
||||
|
||||
func StopSocketServer() {
|
||||
if IsClientConnected() {
|
||||
client.Close()
|
||||
}
|
||||
if server != nil {
|
||||
server.Close()
|
||||
}
|
||||
}
|
||||
-187
@@ -1,187 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
const (
|
||||
DbDriver = "sqlite3"
|
||||
TableCreate = `
|
||||
create table IF NOT EXISTS telegram (
|
||||
[tele_id] INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
[tele_recv_time] TIMESTAMP NOT NULL DEFAULT (datetime('now', 'localtime')),
|
||||
[tele_processed] int(1) NOT NULL DEFAULT 0,
|
||||
[tele_text] TEXT NOT NULL
|
||||
)
|
||||
`
|
||||
InsertNew = "insert into telegram (tele_text) values (?)"
|
||||
InsertOld = "insert into telegram (tele_text, tele_processed) values (?, 1)"
|
||||
|
||||
count = "select count(*) from telegram where tele_processed=0"
|
||||
load = "select tele_id, tele_text from telegram where tele_processed=0 Limit 100"
|
||||
update = "update telegram set tele_processed = 1 where tele_id=?"
|
||||
)
|
||||
|
||||
type telegram struct {
|
||||
id int64
|
||||
text string
|
||||
}
|
||||
|
||||
var (
|
||||
dbFile string
|
||||
initTable bool
|
||||
initialized bool
|
||||
isWriting bool
|
||||
)
|
||||
|
||||
func getDb() *sql.DB {
|
||||
//log.Println("init database ", initTable)
|
||||
if initTable && !initialized {
|
||||
Log.Info("remove database ", zap.String("file", dbFile))
|
||||
os.Remove(dbFile)
|
||||
}
|
||||
db, err := sql.Open(DbDriver, dbFile)
|
||||
if err != nil {
|
||||
Log.Fatal("error in open database ", zap.Error(err))
|
||||
}
|
||||
return db
|
||||
}
|
||||
|
||||
func InitDb(file string, init bool) error {
|
||||
dbFile = file
|
||||
initTable = init
|
||||
var dbErr error
|
||||
db := getDb()
|
||||
defer db.Close()
|
||||
_, dbErr = db.Exec(TableCreate)
|
||||
checkDbErr(dbErr, "error in create table")
|
||||
Log.Info("database initialized")
|
||||
initialized = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func InsertTelegram(teleString string) error {
|
||||
getWriteLock()
|
||||
db := getDb()
|
||||
defer db.Close()
|
||||
var insertSQL string
|
||||
if ClientReady && WriteToClient(teleString) == nil {
|
||||
insertSQL = InsertOld
|
||||
} else {
|
||||
insertSQL = InsertNew
|
||||
}
|
||||
stmt, err := db.Prepare(insertSQL)
|
||||
if err != nil {
|
||||
isWriting = false
|
||||
Log.Error("error in prepare insert telegram ", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
|
||||
defer stmt.Close()
|
||||
result, err := stmt.Exec(teleString)
|
||||
if err != nil {
|
||||
isWriting = false
|
||||
Log.Error("error in insert telegram ", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
id, _ := result.LastInsertId()
|
||||
isWriting = false
|
||||
if insertSQL == InsertOld {
|
||||
Log.Info("telegram processed ", zap.Int64("id", id))
|
||||
} else {
|
||||
Log.Info("telegram saved ", zap.Int64("id", id))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func LoadUnprocessed() {
|
||||
Log.Info("loading telegram")
|
||||
db := getDb()
|
||||
defer db.Close()
|
||||
for {
|
||||
telegrams := getTelegram(db)
|
||||
if len(telegrams) < 1 {
|
||||
break
|
||||
}
|
||||
for _, t := range telegrams {
|
||||
if WriteToClient(t.text) == nil {
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
processed(db, t)
|
||||
}
|
||||
}
|
||||
}
|
||||
ClientReady = true
|
||||
Log.Info("client is ready for new telegram")
|
||||
}
|
||||
|
||||
func getTelegram(db *sql.DB) []telegram {
|
||||
var telegrams []telegram
|
||||
num := countTelegram()
|
||||
Log.Info("telegram ", zap.Int64("unprocessed", num))
|
||||
if num < 1 {
|
||||
return telegrams
|
||||
}
|
||||
Log.Info("try to load unprocessed telegram ")
|
||||
rows, err := db.Query(load)
|
||||
checkDbErr(err, "error in query telegram")
|
||||
defer rows.Close()
|
||||
for rows.Next() {
|
||||
var t telegram
|
||||
err := rows.Scan(&t.id, &t.text)
|
||||
checkDbErr(err, "error in scan telegram")
|
||||
telegrams = append(telegrams, t)
|
||||
}
|
||||
return telegrams
|
||||
}
|
||||
|
||||
func processed(db *sql.DB, t telegram) {
|
||||
getWriteLock()
|
||||
tx, err := db.Begin()
|
||||
checkDbErr(err, "error in begin update")
|
||||
stmt, err := db.Prepare(update)
|
||||
defer stmt.Close()
|
||||
checkDbErr(err, "error in create update statement")
|
||||
result, err := stmt.Exec(t.id)
|
||||
checkDbErr(err, "error in update telegram status")
|
||||
rows, _ := result.RowsAffected()
|
||||
Log.Info("telegram update ", zap.Int64("id", t.id), zap.Bool("result", rows > 0))
|
||||
isWriting = false
|
||||
Log.Debug("release lock")
|
||||
tx.Commit()
|
||||
}
|
||||
|
||||
func countTelegram() int64 {
|
||||
db := getDb()
|
||||
var countErr error
|
||||
stmt, _ := db.Prepare(count)
|
||||
defer stmt.Close()
|
||||
var num int64
|
||||
countErr = stmt.QueryRow().Scan(&num)
|
||||
checkDbErr(countErr, "error in count telegram")
|
||||
return num
|
||||
}
|
||||
|
||||
func getWriteLock() {
|
||||
count := 0
|
||||
Log.Debug("with ", zap.Bool("lock", isWriting))
|
||||
for isWriting && count < 10 {
|
||||
Log.Info("waiting for write ")
|
||||
time.Sleep(5 * time.Second)
|
||||
count++
|
||||
}
|
||||
if count >= 10 {
|
||||
Log.Fatal("database locked")
|
||||
}
|
||||
isWriting = true
|
||||
}
|
||||
|
||||
func checkDbErr(err error, msg string) {
|
||||
if err != nil {
|
||||
Log.Fatal(msg, zap.Error(err))
|
||||
}
|
||||
}
|
||||
@@ -1,63 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"go.uber.org/zap"
|
||||
"regexp"
|
||||
"strings"
|
||||
|
||||
rotatelogs "github.com/lestrrat-go/file-rotatelogs"
|
||||
)
|
||||
|
||||
const (
|
||||
EndTag = "NNNN"
|
||||
Expression = "(?s)ZCZC.*?NNNN"
|
||||
MaxBufferSize = 65536
|
||||
BufferTrimSize = 32768
|
||||
)
|
||||
|
||||
var buffer string
|
||||
var exp, _ = regexp.Compile(Expression)
|
||||
var RawLog *rotatelogs.RotateLogs
|
||||
var Tcp bool
|
||||
var Pulsar bool
|
||||
|
||||
func Append(data string) bool {
|
||||
buffer += data + "\n"
|
||||
if len(buffer) > MaxBufferSize {
|
||||
Log.Warn("buffer overflow, truncating",
|
||||
zap.Int("size", len(buffer)),
|
||||
zap.Int("max", MaxBufferSize))
|
||||
buffer = buffer[len(buffer)-BufferTrimSize:]
|
||||
}
|
||||
return check()
|
||||
}
|
||||
|
||||
func check() bool {
|
||||
teleString := buffer
|
||||
if strings.Contains(teleString, EndTag) {
|
||||
loc := exp.FindStringIndex(teleString)
|
||||
if loc != nil {
|
||||
telegram := teleString[loc[0]:loc[1]]
|
||||
telegram = removeEmpty(telegram) + "\n\n"
|
||||
telegram += "\n"
|
||||
_, _ = RawLog.Write([]byte(telegram + "\n"))
|
||||
buffer = teleString[loc[1]:]
|
||||
|
||||
// 先持久化,再发送:即使发送失败,数据已安全落盘
|
||||
if err := InsertTelegram(telegram); err != nil {
|
||||
Log.Error("error in save telegram to db", zap.Error(err))
|
||||
}
|
||||
if Pulsar {
|
||||
if err := PulsarSend(telegram); err != nil {
|
||||
Log.Error("error in send to pulsar", zap.Error(err))
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func removeEmpty(s string) string {
|
||||
return regexp.MustCompile(`[\t\r\n]+`).ReplaceAllString(strings.TrimSpace(s), "\n")
|
||||
}
|
||||
Reference in New Issue
Block a user