Compare commits
10
Commits
980c52c940
...
4191868f00
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4191868f00 | ||
|
|
46be9f483f | ||
|
|
7f65004cd4 | ||
|
|
86d473c03d | ||
|
|
419f25f0dd | ||
|
|
5f3334ec87 | ||
|
|
133c80088c | ||
|
|
fe1a8169e6 | ||
|
|
9ec63dd8b8 | ||
|
|
b9a826faf8 |
@@ -1,64 +1,202 @@
|
||||
# 电报接收
|
||||
# tele-recv 串口电报接收器
|
||||
|
||||
[](http://d2.int.it2000.com.cn/airport/tele-recv)
|
||||
|
||||
机场串口电报数据缓存程序。从串口读取电报,保存,等待处理程序处理
|
||||
机场串口电报数据缓存与分发程序。从串口读取电报(ZCZC...NNNN 格式),持久化到 SQLite,并通过 TCP 和/或 Apache Pulsar 分发。
|
||||
|
||||
## 下载
|
||||
## 架构
|
||||
|
||||
下载发布版本 [release](https://gitea.int.it2000.com.cn/airport/tele-recv/releases)
|
||||
|
||||
## 安装
|
||||
|
||||
解压下载文件 tele-recv-XXX.tag.gz
|
||||
|
||||
如果是linux
|
||||
|
||||
```Shell
|
||||
tar xzvf tele-recv-XXX.tar.gz
|
||||
```
|
||||
串口 (ttyS0) → serial.Port → telegram.Parser → storage.Store (SQLite)
|
||||
↓
|
||||
transport.Sender
|
||||
↓
|
||||
┌──────────┴──────────┐
|
||||
↓ ↓
|
||||
TCP Client Apache Pulsar
|
||||
```
|
||||
|
||||
文件包中包含两个可执行文件和一个.yaml的配置文件
|
||||
核心原则:**先持久化,再分发**。无论配置何种传输模式,电报始终先写入 SQLite,确保数据不丢失。
|
||||
|
||||
tele-recv-linux: linux 可执行文件
|
||||
tele-recv-win64.exe: windows 64位可执行文件
|
||||
## 项目结构
|
||||
|
||||
### 配置
|
||||
```
|
||||
tele-recv/
|
||||
├── app/ # 应用生命周期管理 (context.Context + 信号处理)
|
||||
├── cmd/ # CLI 入口 (cobra: root, start, test)
|
||||
├── config/ # 配置结构体与 viper 加载
|
||||
├── serial/ # 串口读接口与实现 (Reader 接口)
|
||||
├── telegram/ # 电报解析器 (ZCZC...NNNN 提取)
|
||||
├── storage/ # SQLite 持久化 (Repository 接口)
|
||||
├── transport/ # 传输层 (Sender 接口: TCP, Pulsar, MultiSender)
|
||||
├── main.go # 程序入口
|
||||
├── Taskfile.yml # 构建/测试任务定义
|
||||
└── telegram.yaml # 配置文件
|
||||
```
|
||||
|
||||
| 选项 | 含义 | 取值 | 例子 |
|
||||
|------- |---------------------|-----------------|------------------------|
|
||||
|device | 串口设备名称 | ttyS1 | win: COM1 linux: ttyS1 |
|
||||
|baudrate| 波特率 | 9600 | 19200,38400 .... |
|
||||
|lograw | 是否向控制台输出通讯内容| ture | true/false |
|
||||
|file | sqlite 文件名 | telegram.db | tele.db |
|
||||
|init | 是否新建电报数据库 | true | true/false |
|
||||
|address | 本地处理监听端口 | 127.0.0.1:6000 | |
|
||||
## 技术栈
|
||||
|
||||
| 组件 | 技术 | 说明 |
|
||||
|------|------|------|
|
||||
| 语言 | Go 1.21+ | 需 Go 1.21 或更高版本 |
|
||||
| CLI | cobra + viper | 命令解析与配置管理 |
|
||||
| 串口 | github.com/argandas/serial | 串口通信 |
|
||||
| 存储 | github.com/mattn/go-sqlite3 | SQLite 驱动 (需要 CGO) |
|
||||
| 消息 | github.com/apache/pulsar-client-go | Apache Pulsar 集成 |
|
||||
| 日志 | go.uber.org/zap | 结构化日志 |
|
||||
| 轮转 | github.com/lestrrat-go/file-rotatelogs | 日志文件轮转 |
|
||||
|
||||
## 快速开始
|
||||
|
||||
### 前置条件
|
||||
|
||||
- Go 1.21+
|
||||
- 如需交叉编译:`x86_64-linux-gnu-gcc` (Linux) 或 `x86_64-w64-mingw32-gcc` (Windows)
|
||||
|
||||
### 安装 Task (推荐)
|
||||
|
||||
```bash
|
||||
go install github.com/go-task/task/v3/cmd/task@latest
|
||||
```
|
||||
|
||||
### 常用命令
|
||||
|
||||
```bash
|
||||
task deps # 安装依赖
|
||||
task build # 构建当前平台二进制
|
||||
task build-all # 构建 Linux + Windows 二进制
|
||||
task run # 构建并启动 tele-recv start
|
||||
task run-test # 构建并启动 tele-recv test
|
||||
task test # 运行所有测试
|
||||
task test-race # 运行竞态检测测试
|
||||
task test-cover # 运行测试并生成覆盖率报告
|
||||
task lint # go vet 静态检查
|
||||
task emu # 创建虚拟串口对 (ttyS0 ↔ ttyS1)
|
||||
task clean # 清理构建产物
|
||||
task dist # 构建全平台并打包 tar.gz
|
||||
```
|
||||
|
||||
### 手动构建
|
||||
|
||||
```bash
|
||||
go mod tidy
|
||||
go build -o tele-recv .
|
||||
```
|
||||
|
||||
### 交叉编译
|
||||
|
||||
```bash
|
||||
# Linux (需要 CGO 和 x86_64-linux-gnu-gcc)
|
||||
CGO_ENABLED=1 GOOS=linux GOARCH=amd64 CC=x86_64-linux-gnu-gcc go build -o tele-recv-linux .
|
||||
|
||||
# Windows (需要 CGO 和 x86_64-w64-mingw32-gcc)
|
||||
CGO_ENABLED=1 GOOS=windows GOARCH=amd64 CC=x86_64-w64-mingw32-gcc go build -o tele-recv-win64.exe .
|
||||
```
|
||||
|
||||
### 模拟串口
|
||||
|
||||
```bash
|
||||
task emu
|
||||
# 创建 /tmp/ttyS0 ↔ /tmp/ttyS1 虚拟串口对
|
||||
```
|
||||
|
||||
## 配置
|
||||
|
||||
编辑 `telegram.yaml`:
|
||||
|
||||
```yaml
|
||||
serial:
|
||||
device: /tmp/ttyS1 # 串口设备
|
||||
baudrate: 9600 # 波特率
|
||||
lograw: true # 控制台输出原始数据
|
||||
|
||||
telegram:
|
||||
tcp: false # 启用 TCP 分发
|
||||
pulsar: true # 启用 Pulsar 分发
|
||||
|
||||
sqlite:
|
||||
file: telegram.db # 数据库文件
|
||||
init: false # 启动时重建数据库
|
||||
|
||||
socket:
|
||||
address: 127.0.0.1:6000 # TCP 监听地址
|
||||
|
||||
pulsar:
|
||||
url: pulsar://localhost:6650
|
||||
topic: telegram-raw
|
||||
name: serial-reader
|
||||
|
||||
log:
|
||||
dir: ./logs # 日志目录
|
||||
maxage: 60 # 日志保留天数
|
||||
rotatehour: 1 # 日志轮转间隔(小时)
|
||||
```
|
||||
|
||||
### 配置选项
|
||||
|
||||
| 选项 | 含义 | 默认值 |
|
||||
|------|------|--------|
|
||||
| `serial.device` | 串口设备名 | — |
|
||||
| `serial.baudrate` | 波特率 | — |
|
||||
| `serial.lograw` | 控制台输出原始数据 | `false` |
|
||||
| `telegram.tcp` | 启用 TCP 分发 | `false` |
|
||||
| `telegram.pulsar` | 启用 Pulsar 分发 | `false` |
|
||||
| `sqlite.file` | SQLite 数据库文件 | — |
|
||||
| `sqlite.init` | 启动时重建数据库 | `false` |
|
||||
| `socket.address` | TCP 监听地址 | — |
|
||||
| `pulsar.url` | Pulsar 代理地址 | — |
|
||||
| `pulsar.topic` | Pulsar 主题 | — |
|
||||
| `pulsar.name` | Pulsar 生产者名称 | — |
|
||||
| `log.dir` | 日志目录 | `./logs` |
|
||||
| `log.maxage` | 日志保留天数 | `60` |
|
||||
| `log.rotatehour` | 日志轮转间隔(小时) | `1` |
|
||||
|
||||
> 注意:TCP 和 Pulsar 可以同时启用。电报会先写入 SQLite,再通过所有启用的传输通道分发。
|
||||
|
||||
## 运行
|
||||
|
||||
### 测试环境
|
||||
|
||||
linux:
|
||||
```Shell
|
||||
./tele-recv-linux test
|
||||
检查串口和数据库是否可用:
|
||||
|
||||
```bash
|
||||
./tele-recv test
|
||||
```
|
||||
|
||||
windows:
|
||||
```Shell
|
||||
tele-recv-win64.exe test
|
||||
### 启动服务
|
||||
|
||||
```bash
|
||||
./tele-recv start
|
||||
```
|
||||
|
||||
测试运行的环境是否能正确开始串口设备,以及初始化数据库
|
||||
服务启动后将:
|
||||
1. 打开串口设备
|
||||
2. 从串口读取数据,解析 ZCZC...NNNN 电报
|
||||
3. 写入 SQLite 持久化
|
||||
4. 通过 TCP 和/或 Pulsar 分发
|
||||
|
||||
### 开始接收电报
|
||||
## 测试
|
||||
|
||||
linux:
|
||||
```Shell
|
||||
./tele-recv-linux start
|
||||
```bash
|
||||
# 运行所有测试
|
||||
task test
|
||||
|
||||
# 带竞态检测
|
||||
task test-race
|
||||
|
||||
# 覆盖率报告
|
||||
task test-cover
|
||||
|
||||
# 手动运行
|
||||
go test -race -v ./...
|
||||
```
|
||||
|
||||
windows:
|
||||
```Shell
|
||||
tele-recv-win64.exe start
|
||||
```
|
||||
当前测试覆盖 6 个包,30+ 测试用例,`go test -race` 零竞态。
|
||||
|
||||
开始从串口设备读取数据,并开启TCP服务等待处理程序连接。
|
||||
如果没有处理程序,电报则保存在本地数据库中。
|
||||
等到处理程序连上端口,则把未处理电报以此发送给处理程序。
|
||||
## 下载
|
||||
|
||||
发布版本:[release](https://gitea.int.it2000.com.cn/airport/tele-recv/releases)
|
||||
|
||||
## 许可证
|
||||
|
||||
详见 [LICENSE](LICENSE)
|
||||
+150
@@ -0,0 +1,150 @@
|
||||
# https://taskfile.dev
|
||||
version: "3"
|
||||
|
||||
vars:
|
||||
APP: tele-recv
|
||||
BIN_DIR: bin
|
||||
CONFIG: telegram.yaml
|
||||
GO_VERSION: "1.15"
|
||||
LDFLAGS: '-s -w'
|
||||
BUILD_DIR: "{{.BIN_DIR}}/{{.APP}}"
|
||||
|
||||
tasks:
|
||||
default:
|
||||
desc: Show available tasks
|
||||
cmd: task --list-all
|
||||
|
||||
# ─── Build ────────────────────────────────────────────────────────────────────
|
||||
|
||||
build:
|
||||
desc: Build the binary for the current platform
|
||||
deps: [deps]
|
||||
cmds:
|
||||
- mkdir -p "{{.BIN_DIR}}"
|
||||
- CGO_ENABLED=1 go build -ldflags="{{.LDFLAGS}}" -o "{{.BIN_DIR}}/{{.APP}}" .
|
||||
sources:
|
||||
- "**/*.go"
|
||||
- go.mod
|
||||
- go.sum
|
||||
generates:
|
||||
- "{{.BIN_DIR}}/{{.APP}}"
|
||||
|
||||
build-all:
|
||||
desc: Cross-compile for linux/amd64 and windows/amd64
|
||||
cmds:
|
||||
- task: build-linux
|
||||
- task: build-windows
|
||||
|
||||
build-linux:
|
||||
desc: Build for linux/amd64 (requires CGO cross-compiler)
|
||||
env:
|
||||
GOOS: linux
|
||||
GOARCH: amd64
|
||||
CGO_ENABLED: "1"
|
||||
CC: x86_64-linux-gnu-gcc
|
||||
cmds:
|
||||
- mkdir -p "{{.BIN_DIR}}"
|
||||
- go build -ldflags="{{.LDFLAGS}}" -o "{{.BIN_DIR}}/{{.APP}}-linux" .
|
||||
sources:
|
||||
- "**/*.go"
|
||||
- go.mod
|
||||
- go.sum
|
||||
generates:
|
||||
- "{{.BIN_DIR}}/{{.APP}}-linux"
|
||||
silent: false
|
||||
|
||||
build-windows:
|
||||
desc: Build for windows/amd64 (requires MinGW cross-compiler)
|
||||
env:
|
||||
GOOS: windows
|
||||
GOARCH: amd64
|
||||
CGO_ENABLED: "1"
|
||||
CC: x86_64-w64-mingw32-gcc
|
||||
cmds:
|
||||
- mkdir -p "{{.BIN_DIR}}"
|
||||
- go build -ldflags="{{.LDFLAGS}}" -o "{{.BIN_DIR}}/{{.APP}}-win64.exe" .
|
||||
sources:
|
||||
- "**/*.go"
|
||||
- go.mod
|
||||
- go.sum
|
||||
generates:
|
||||
- "{{.BIN_DIR}}/{{.APP}}-win64.exe"
|
||||
silent: false
|
||||
|
||||
# ─── Test ─────────────────────────────────────────────────────────────────────
|
||||
|
||||
test:
|
||||
desc: Run all tests
|
||||
cmds:
|
||||
- go test ./...
|
||||
|
||||
test-race:
|
||||
desc: Run tests with race detector
|
||||
cmds:
|
||||
- go test -race ./...
|
||||
|
||||
test-cover:
|
||||
desc: Run tests with coverage report
|
||||
cmds:
|
||||
- go test -coverprofile=coverage.out ./...
|
||||
- go tool cover -func=coverage.out
|
||||
|
||||
# ─── Run ──────────────────────────────────────────────────────────────────────
|
||||
|
||||
run:
|
||||
desc: Run the service (start subcommand)
|
||||
deps: [build]
|
||||
cmds:
|
||||
- "{{.BIN_DIR}}/{{.APP}} start"
|
||||
|
||||
run-test:
|
||||
desc: Run the environment test (test subcommand)
|
||||
deps: [build]
|
||||
cmds:
|
||||
- "{{.BIN_DIR}}/{{.APP}} test"
|
||||
|
||||
# ─── Development helpers ──────────────────────────────────────────────────────
|
||||
|
||||
emu:
|
||||
desc: Create a virtual serial port pair (ttyS0 ↔ ttyS1) with socat
|
||||
cmds:
|
||||
- socat PTY,link=/tmp/ttyS0 PTY,link=/tmp/ttyS1
|
||||
silent: false
|
||||
|
||||
deps:
|
||||
desc: Tidy and download Go module dependencies
|
||||
cmds:
|
||||
- go mod tidy
|
||||
- go mod download
|
||||
sources:
|
||||
- go.mod
|
||||
- go.sum
|
||||
generates:
|
||||
- go.sum
|
||||
|
||||
lint:
|
||||
desc: Run go vet on all packages
|
||||
cmds:
|
||||
- go vet ./...
|
||||
|
||||
clean:
|
||||
desc: Remove build artifacts, logs, and database
|
||||
cmds:
|
||||
- rm -rf "{{.BIN_DIR}}"
|
||||
- rm -f telegram.db
|
||||
- rm -rf ./logs/
|
||||
silent: false
|
||||
|
||||
# ─── Distribution package ─────────────────────────────────────────────────────
|
||||
|
||||
dist:
|
||||
desc: Build all platforms and package into a tarball
|
||||
cmds:
|
||||
- task: build-all
|
||||
- mkdir -p "{{.BUILD_DIR}}"
|
||||
- cp "{{.BIN_DIR}}/{{.APP}}-linux" "{{.BUILD_DIR}}/"
|
||||
- cp "{{.BIN_DIR}}/{{.APP}}-win64.exe" "{{.BUILD_DIR}}/"
|
||||
- cp "{{.CONFIG}}" "{{.BUILD_DIR}}/"
|
||||
- cd "{{.BIN_DIR}}" && tar czvf "{{.APP}}.tar.gz" "{{.APP}}/" && mv "{{.APP}}.tar.gz" .
|
||||
- rm -rf "{{.BUILD_DIR}}"
|
||||
silent: false
|
||||
+268
@@ -0,0 +1,268 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
rotatelogs "github.com/lestrrat-go/file-rotatelogs"
|
||||
|
||||
"it2000.com.cn/tele-recv/config"
|
||||
"it2000.com.cn/tele-recv/serial"
|
||||
"it2000.com.cn/tele-recv/telegram"
|
||||
"it2000.com.cn/tele-recv/storage"
|
||||
"it2000.com.cn/tele-recv/transport"
|
||||
)
|
||||
|
||||
// App orchestrates the complete telegram receive pipeline.
|
||||
type App struct {
|
||||
cfg *config.Config
|
||||
logger *zap.Logger
|
||||
rawLog *rotatelogs.RotateLogs
|
||||
port *serial.Port
|
||||
parser *telegram.Parser
|
||||
store storage.Repository
|
||||
sender transport.Sender
|
||||
tcpSrv *transport.TCPServer
|
||||
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// New creates a new App from configuration.
|
||||
func New(cfg *config.Config) (*App, error) {
|
||||
// Initialize logger
|
||||
logger, rawLog, err := initLoggers(cfg)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("init log: %w", err)
|
||||
}
|
||||
|
||||
// Initialize parser
|
||||
parser := telegram.New(rawLog)
|
||||
|
||||
// Initialize store
|
||||
store, err := storage.New(cfg.SQLite.File, cfg.SQLite.Init)
|
||||
if err != nil {
|
||||
logger.Error("failed to init store", zap.Error(err))
|
||||
// Non-fatal: we can still run without persistence
|
||||
store = nil
|
||||
}
|
||||
|
||||
// Initialize sender
|
||||
var sender transport.Sender
|
||||
var tcpSrv *transport.TCPServer
|
||||
|
||||
if cfg.Telegram.TCP {
|
||||
tcpSender := transport.NewTCPSender(cfg.Socket.Address)
|
||||
tcpSrv = transport.NewTCPServer(cfg.Socket.Address, tcpSender, nil)
|
||||
sender = tcpSender
|
||||
}
|
||||
|
||||
if cfg.Telegram.Pulsar {
|
||||
pulsarSender, err := transport.NewPulsarSender(
|
||||
cfg.Pulsar.URL, cfg.Pulsar.Topic, cfg.Pulsar.Name,
|
||||
)
|
||||
if err != nil {
|
||||
logger.Warn("failed to init pulsar sender", zap.Error(err))
|
||||
} else {
|
||||
if sender != nil {
|
||||
sender = transport.NewMultiSender(sender, pulsarSender)
|
||||
} else {
|
||||
sender = pulsarSender
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
return &App{
|
||||
cfg: cfg,
|
||||
logger: logger,
|
||||
rawLog: rawLog,
|
||||
parser: parser,
|
||||
store: store,
|
||||
sender: sender,
|
||||
tcpSrv: tcpSrv,
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Run starts the application and blocks until a signal or error.
|
||||
func (a *App) Run() error {
|
||||
defer a.logger.Sync()
|
||||
defer a.cleanup()
|
||||
|
||||
a.logger.Info("starting telegram receiver",
|
||||
zap.String("device", a.cfg.Serial.Device),
|
||||
zap.Int("baudrate", a.cfg.Serial.Baudrate))
|
||||
|
||||
// Start TCP server if configured
|
||||
if a.tcpSrv != nil {
|
||||
a.wg.Add(1)
|
||||
go func() {
|
||||
defer a.wg.Done()
|
||||
if err := a.tcpSrv.Run(a.ctx); err != nil {
|
||||
a.logger.Error("tcp server error", zap.Error(err))
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// Signal handling
|
||||
sigCh := make(chan os.Signal, 1)
|
||||
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
|
||||
|
||||
// Main loop: read from serial port, parse, store, send
|
||||
const maxRetries = 3
|
||||
runCtx, runCancel := context.WithCancel(a.ctx)
|
||||
defer runCancel()
|
||||
|
||||
go func() {
|
||||
<-sigCh
|
||||
a.logger.Info("received signal, shutting down")
|
||||
runCancel()
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-runCtx.Done():
|
||||
return nil
|
||||
default:
|
||||
}
|
||||
|
||||
if a.port == nil || !a.port.IsOpen() {
|
||||
if err := a.openPort(maxRetries); err != nil {
|
||||
if err == context.Canceled {
|
||||
return nil
|
||||
}
|
||||
a.logger.Fatal("failed to open serial port", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
line, err := a.port.ReadLine()
|
||||
if err != nil {
|
||||
a.logger.Warn("serial read error", zap.Error(err))
|
||||
a.port.Close()
|
||||
continue
|
||||
}
|
||||
|
||||
if a.cfg.Serial.LogRaw && len(line) > 0 {
|
||||
fmt.Println(line)
|
||||
}
|
||||
|
||||
telegrams := a.parser.Append(line)
|
||||
for _, telegram := range telegrams {
|
||||
a.dispatch(telegram)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (a *App) openPort(maxRetries int) error {
|
||||
a.logger.Info("opening serial port",
|
||||
zap.String("device", a.cfg.Serial.Device),
|
||||
zap.Int("baudrate", a.cfg.Serial.Baudrate))
|
||||
|
||||
var lastErr error
|
||||
for i := 1; i <= maxRetries; i++ {
|
||||
port, err := serial.Open(a.cfg.Serial.Device, a.cfg.Serial.Baudrate)
|
||||
if err == nil {
|
||||
a.port = port
|
||||
return nil
|
||||
}
|
||||
lastErr = err
|
||||
a.logger.Warn("serial port open failed, retrying",
|
||||
zap.Int("attempt", i),
|
||||
zap.Error(err))
|
||||
|
||||
select {
|
||||
case <-a.ctx.Done():
|
||||
return context.Canceled
|
||||
case <-time.After(time.Duration(1<<(i-1)) * time.Second):
|
||||
}
|
||||
}
|
||||
return lastErr
|
||||
}
|
||||
|
||||
func (a *App) dispatch(telegram string) {
|
||||
// Always persist first
|
||||
if a.store != nil {
|
||||
if err := a.store.Insert(telegram); err != nil {
|
||||
a.logger.Error("failed to persist telegram", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
// Then send to transport(s)
|
||||
if a.sender != nil {
|
||||
if err := a.sender.Send(telegram); err != nil {
|
||||
a.logger.Error("failed to send telegram", zap.Error(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (a *App) cleanup() {
|
||||
a.logger.Info("shutting down")
|
||||
|
||||
if a.tcpSrv != nil {
|
||||
a.tcpSrv.Stop()
|
||||
}
|
||||
|
||||
if a.port != nil {
|
||||
a.port.Close()
|
||||
}
|
||||
|
||||
if a.sender != nil {
|
||||
a.sender.Close()
|
||||
}
|
||||
|
||||
if a.store != nil {
|
||||
a.store.Close()
|
||||
}
|
||||
|
||||
a.wg.Wait()
|
||||
}
|
||||
|
||||
func initLoggers(cfg *config.Config) (*zap.Logger, *rotatelogs.RotateLogs, error) {
|
||||
logFile := cfg.Log.Dir + "/telegram-%Y-%m-%d-%H.log"
|
||||
rotator, err := rotatelogs.New(
|
||||
logFile,
|
||||
rotatelogs.WithMaxAge(time.Duration(cfg.Log.MaxAge)*24*time.Hour),
|
||||
rotatelogs.WithRotationTime(time.Duration(cfg.Log.RotateHour)*time.Hour),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
rawFile := cfg.Log.Dir + "/raw-%Y-%m-%d-%H.txt"
|
||||
rawLog, err := rotatelogs.New(
|
||||
rawFile,
|
||||
rotatelogs.WithMaxAge(time.Duration(cfg.Log.MaxAge)*24*time.Hour),
|
||||
rotatelogs.WithRotationTime(time.Duration(cfg.Log.RotateHour)*time.Hour),
|
||||
)
|
||||
if err != nil {
|
||||
rotator.Close()
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
filePriority := zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
|
||||
return lvl >= zapcore.DebugLevel
|
||||
})
|
||||
stdoutPriority := zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
|
||||
return lvl >= zapcore.InfoLevel
|
||||
})
|
||||
|
||||
encoder := zapcore.NewConsoleEncoder(zap.NewDevelopmentEncoderConfig())
|
||||
core := zapcore.NewTee(
|
||||
zapcore.NewCore(encoder, zapcore.Lock(os.Stdout), stdoutPriority),
|
||||
zapcore.NewCore(encoder, zapcore.AddSync(rotator), filePriority),
|
||||
)
|
||||
|
||||
logger := zap.New(core)
|
||||
|
||||
return logger, rawLog, nil
|
||||
}
|
||||
+182
@@ -0,0 +1,182 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"it2000.com.cn/tele-recv/config"
|
||||
"it2000.com.cn/tele-recv/storage"
|
||||
)
|
||||
|
||||
// mockSender implements transport.Sender for testing.
|
||||
type mockSender struct {
|
||||
mu sync.Mutex
|
||||
sent []string
|
||||
sendErr error
|
||||
closeErr error
|
||||
}
|
||||
|
||||
func (m *mockSender) Send(telegram string) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if m.sendErr != nil {
|
||||
return m.sendErr
|
||||
}
|
||||
m.sent = append(m.sent, telegram)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockSender) Close() error { return m.closeErr }
|
||||
|
||||
// mockStore implements storage insertion for testing.
|
||||
type mockStore struct {
|
||||
mu sync.Mutex
|
||||
telegrams []string
|
||||
insertErr error
|
||||
}
|
||||
|
||||
func (m *mockStore) Insert(telegram string) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if m.insertErr != nil {
|
||||
return m.insertErr
|
||||
}
|
||||
m.telegrams = append(m.telegrams, telegram)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockStore) LoadUnprocessed() ([]storage.Telegram, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (m *mockStore) MarkProcessed(id int64) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockStore) Close() error { return nil }
|
||||
|
||||
func TestNewApp_NilConfig(t *testing.T) {
|
||||
// Should not panic, but serial port open will fail later
|
||||
cfg := &config.Config{
|
||||
Serial: config.SerialConfig{
|
||||
Device: "/nonexistent",
|
||||
Baudrate: 9600,
|
||||
},
|
||||
SQLite: config.SQLiteConfig{
|
||||
File: ":memory:",
|
||||
Init: true,
|
||||
},
|
||||
Log: config.LogConfig{
|
||||
Dir: t.TempDir(),
|
||||
MaxAge: 1,
|
||||
RotateHour: 24,
|
||||
},
|
||||
}
|
||||
|
||||
a, err := New(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
if a == nil {
|
||||
t.Fatal("expected non-nil app")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatch_StoreAndSend(t *testing.T) {
|
||||
store := &mockStore{}
|
||||
sender := &mockSender{}
|
||||
cfg := &config.Config{
|
||||
SQLite: config.SQLiteConfig{File: ":memory:"},
|
||||
Log: config.LogConfig{Dir: t.TempDir(), MaxAge: 1, RotateHour: 24},
|
||||
}
|
||||
|
||||
app, err := New(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
app.store = store
|
||||
app.sender = sender
|
||||
|
||||
app.dispatch("ZCZC TEST NNNN")
|
||||
|
||||
if len(store.telegrams) != 1 {
|
||||
t.Errorf("expected 1 stored telegram, got %d", len(store.telegrams))
|
||||
}
|
||||
if len(sender.sent) != 1 {
|
||||
t.Errorf("expected 1 sent telegram, got %d", len(sender.sent))
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatch_StoreError(t *testing.T) {
|
||||
store := &mockStore{insertErr: errors.New("db error")}
|
||||
sender := &mockSender{}
|
||||
cfg := &config.Config{
|
||||
SQLite: config.SQLiteConfig{File: ":memory:"},
|
||||
Log: config.LogConfig{Dir: t.TempDir(), MaxAge: 1, RotateHour: 24},
|
||||
}
|
||||
|
||||
app, err := New(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
app.store = store
|
||||
app.sender = sender
|
||||
|
||||
// Should not panic on store error, should still send
|
||||
app.dispatch("ZCZC TEST NNNN")
|
||||
|
||||
if len(sender.sent) != 1 {
|
||||
t.Errorf("expected 1 sent telegram despite store error, got %d", len(sender.sent))
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatch_SendError(t *testing.T) {
|
||||
store := &mockStore{}
|
||||
sender := &mockSender{sendErr: errors.New("send error")}
|
||||
cfg := &config.Config{
|
||||
SQLite: config.SQLiteConfig{File: ":memory:"},
|
||||
Log: config.LogConfig{Dir: t.TempDir(), MaxAge: 1, RotateHour: 24},
|
||||
}
|
||||
|
||||
app, err := New(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
app.store = store
|
||||
app.sender = sender
|
||||
|
||||
// Should not panic on send error, should still store
|
||||
app.dispatch("ZCZC TEST NNNN")
|
||||
|
||||
if len(store.telegrams) != 1 {
|
||||
t.Errorf("expected 1 stored telegram despite send error, got %d", len(store.telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRun_CancelledContext(t *testing.T) {
|
||||
cfg := &config.Config{
|
||||
Serial: config.SerialConfig{
|
||||
Device: "/nonexistent",
|
||||
Baudrate: 9600,
|
||||
},
|
||||
SQLite: config.SQLiteConfig{File: ":memory:"},
|
||||
Log: config.LogConfig{Dir: t.TempDir(), MaxAge: 1, RotateHour: 24},
|
||||
}
|
||||
app, err := New(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
|
||||
// Cancel after a short delay so Run() tries opening port and exits cleanly
|
||||
go func() {
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
app.cancel()
|
||||
}()
|
||||
|
||||
err = app.Run()
|
||||
if err != nil {
|
||||
t.Fatalf("Run() should return nil on cancellation, got: %v", err)
|
||||
}
|
||||
}
|
||||
+33
-70
@@ -25,43 +25,38 @@ import (
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
|
||||
"github.com/spf13/viper"
|
||||
"it2000.com.cn/tele-recv/utils"
|
||||
"it2000.com.cn/tele-recv/config"
|
||||
)
|
||||
|
||||
const ModeName = "TELEGRAM_MODE"
|
||||
|
||||
var (
|
||||
cfgFile string
|
||||
device string
|
||||
baudrate int
|
||||
lograw bool
|
||||
|
||||
dbFile string
|
||||
dbInit bool
|
||||
cfgFile string
|
||||
|
||||
// Legacy global vars for backward compatibility with test command
|
||||
device string
|
||||
baudrate int
|
||||
dbFile string
|
||||
dbInit bool
|
||||
socketAddress string
|
||||
pulsarUrl string
|
||||
topic string
|
||||
name string
|
||||
tcp bool
|
||||
pulsar bool
|
||||
lograw bool
|
||||
|
||||
pulsarUrl string
|
||||
topic string
|
||||
name string
|
||||
|
||||
tcp bool
|
||||
pulsar bool
|
||||
Mode string
|
||||
|
||||
// rootCmd represents the base command when called without any subcommands
|
||||
rootCmd = &cobra.Command{
|
||||
Use: "tele-recv",
|
||||
Short: "telegram receiver",
|
||||
Long: `A serial port telegram receiver.
|
||||
Read the telegram from `,
|
||||
// Uncomment the following line if your bare application
|
||||
// has an action associated with it:
|
||||
// Run: func(cmd *cobra.Command, args []string) { },
|
||||
Read the telegram from serial port`,
|
||||
}
|
||||
|
||||
Mode = os.Getenv(ModeName)
|
||||
|
||||
appCfg *config.Config
|
||||
logger *zap.Logger
|
||||
)
|
||||
|
||||
@@ -77,55 +72,33 @@ func Execute() {
|
||||
func init() {
|
||||
cobra.OnInitialize(initConfig)
|
||||
|
||||
// Here you will define your flags and configuration settings.
|
||||
// Cobra supports persistent flags, which, if defined here,
|
||||
// will be global for your application.
|
||||
|
||||
rootCmd.PersistentFlags().StringVar(&cfgFile, "config", "", "config file (telegram.yaml)")
|
||||
// Cobra also supports local flags, which will only run
|
||||
// when this action is called directly.
|
||||
rootCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
|
||||
|
||||
initLog()
|
||||
defer logger.Sync()
|
||||
// l = logger.Sugar()
|
||||
utils.Log = logger
|
||||
|
||||
}
|
||||
|
||||
// initConfig reads in config file and ENV variables if set.
|
||||
func initConfig() {
|
||||
if cfgFile != "" {
|
||||
// Use config file from the flag.
|
||||
viper.SetConfigFile(cfgFile)
|
||||
} else {
|
||||
// Search config in home directory with name ".tele-recv" (without extension).
|
||||
viper.AddConfigPath(".")
|
||||
viper.SetConfigName("telegram")
|
||||
}
|
||||
|
||||
viper.AutomaticEnv() // read in environment variables that match
|
||||
|
||||
// If a config file is found, read it in.
|
||||
if err := viper.ReadInConfig(); err == nil {
|
||||
// fmt.Println("Using config file:", viper.ConfigFileUsed())
|
||||
logger.Info("Using config file ", zap.String("path", viper.ConfigFileUsed()))
|
||||
device = viper.GetString("serial.device")
|
||||
baudrate = viper.GetInt("serial.baudrate")
|
||||
lograw = viper.GetBool("serial.lograw")
|
||||
dbFile = viper.GetString("sqlite.file")
|
||||
dbInit = viper.GetBool("sqlite.init")
|
||||
socketAddress = viper.GetString("socket.address")
|
||||
pulsarUrl = viper.GetString("pulsar.url")
|
||||
topic = viper.GetString("pulsar.topic")
|
||||
name = viper.GetString("pulsar.name")
|
||||
|
||||
tcp = viper.GetBool("telegram.tcp")
|
||||
utils.Tcp = tcp
|
||||
pulsar = viper.GetBool("telegram.pulsar")
|
||||
utils.Pulsar = pulsar
|
||||
cfg, err := config.Load(cfgFile)
|
||||
if err != nil {
|
||||
logger.Warn("config load error, using defaults", zap.Error(err))
|
||||
return
|
||||
}
|
||||
appCfg = cfg
|
||||
|
||||
// Populate legacy globals
|
||||
device = cfg.Serial.Device
|
||||
baudrate = cfg.Serial.Baudrate
|
||||
lograw = cfg.Serial.LogRaw
|
||||
dbFile = cfg.SQLite.File
|
||||
dbInit = cfg.SQLite.Init
|
||||
socketAddress = cfg.Socket.Address
|
||||
pulsarUrl = cfg.Pulsar.URL
|
||||
topic = cfg.Pulsar.Topic
|
||||
name = cfg.Pulsar.Name
|
||||
tcp = cfg.Telegram.TCP
|
||||
pulsar = cfg.Telegram.Pulsar
|
||||
}
|
||||
|
||||
func initLog() {
|
||||
@@ -138,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
|
||||
})
|
||||
@@ -163,5 +127,4 @@ func initLog() {
|
||||
)
|
||||
|
||||
logger = zap.New(core)
|
||||
|
||||
}
|
||||
|
||||
+14
-77
@@ -16,15 +16,10 @@ limitations under the License.
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"it2000.com.cn/tele-recv/utils"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/zap"
|
||||
|
||||
"it2000.com.cn/tele-recv/app"
|
||||
)
|
||||
|
||||
// startCmd represents the start command
|
||||
@@ -33,82 +28,24 @@ var startCmd = &cobra.Command{
|
||||
Short: "start telegram receive service",
|
||||
Long: `open serial port
|
||||
start a tcp server for processing`,
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
start()
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
return start()
|
||||
},
|
||||
}
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(startCmd)
|
||||
|
||||
// Here you will define your flags and configuration settings.
|
||||
|
||||
// Cobra supports Persistent Flags which will work for this command
|
||||
// and all subcommands, e.g.:
|
||||
// startCmd.PersistentFlags().String("foo", "", "A help for foo")
|
||||
|
||||
// Cobra supports local flags which will only run when this command
|
||||
// is called directly, e.g.:
|
||||
// startCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
|
||||
}
|
||||
|
||||
func start() {
|
||||
// Go signal notification works by sending `os.Signal`
|
||||
// values on a channel. We'll create a channel to
|
||||
// receive these notifications (we'll also make one to
|
||||
// notify us when the program can exit).
|
||||
sigs := make(chan os.Signal, 1)
|
||||
done := make(chan bool, 1)
|
||||
|
||||
// `signal.Notify` registers the given channel to
|
||||
// receive notifications of the specified signals.
|
||||
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
|
||||
utils.ServerRunning = true
|
||||
|
||||
// This goroutine executes a blocking receive for
|
||||
// signals. When it gets one it'll print it out
|
||||
// and then notify the program that it can finish.
|
||||
go func() {
|
||||
sig := <-sigs
|
||||
logger.Info("got ", zap.Any("signal", sig))
|
||||
utils.ServerRunning = false
|
||||
if tcp {
|
||||
utils.StopSocketServer()
|
||||
}
|
||||
if pulsar {
|
||||
utils.ClosePulsar()
|
||||
}
|
||||
done <- true
|
||||
}()
|
||||
|
||||
// The program will wait here until it gets the
|
||||
// expected signal (as indicated by the goroutine
|
||||
// above sending a value on `done`) and then exit.
|
||||
logger.Info("awaiting signal")
|
||||
_ = utils.InitDb(dbFile, dbInit)
|
||||
if tcp {
|
||||
go utils.Listen(socketAddress)
|
||||
func start() error {
|
||||
if appCfg == nil {
|
||||
logger.Fatal("configuration not loaded")
|
||||
}
|
||||
for utils.ServerRunning {
|
||||
if !utils.IsPortOpen() {
|
||||
logger.Info("try to open port")
|
||||
err := utils.OpenPort(device, baudrate)
|
||||
if err != nil {
|
||||
logger.Fatal("error in open serial port ", zap.Error(err))
|
||||
}
|
||||
logger.Info("starting read")
|
||||
}
|
||||
buffer, err := utils.ReadPort()
|
||||
if err == nil {
|
||||
if lograw && len(buffer) > 0 {
|
||||
fmt.Println(buffer)
|
||||
}
|
||||
if utils.Append(buffer) {
|
||||
fmt.Println()
|
||||
}
|
||||
}
|
||||
}
|
||||
<-done
|
||||
logger.Info("exiting")
|
||||
|
||||
a, err := app.New(appCfg)
|
||||
if err != nil {
|
||||
logger.Fatal("failed to create app", zap.Error(err))
|
||||
}
|
||||
|
||||
return a.Run()
|
||||
}
|
||||
|
||||
+40
-16
@@ -16,8 +16,14 @@ limitations under the License.
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"net"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"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"
|
||||
)
|
||||
|
||||
// testCmd represents the test command
|
||||
@@ -34,27 +40,45 @@ 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 the Pulsar endpoint without constructing a producer. The legacy Pulsar
|
||||
// client used by the service is not compatible with newer Go runtimes on macOS.
|
||||
endpoint, err := url.Parse(pulsarUrl)
|
||||
if err != nil {
|
||||
logger.Warn("pulsar URL is invalid", zap.Error(err))
|
||||
} else {
|
||||
address := endpoint.Host
|
||||
if endpoint.Port() == "" {
|
||||
address = net.JoinHostPort(endpoint.Hostname(), "6650")
|
||||
}
|
||||
conn, err := net.DialTimeout("tcp", address, 3*time.Second)
|
||||
if err != nil {
|
||||
logger.Warn("pulsar connection failed (expected if no broker)", zap.Error(err))
|
||||
} else {
|
||||
conn.Close()
|
||||
logger.Info("pulsar endpoint ok")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(testCmd)
|
||||
|
||||
// Here you will define your flags and configuration settings.
|
||||
|
||||
// Cobra supports Persistent Flags which will work for this command
|
||||
// and all subcommands, e.g.:
|
||||
// testCmd.PersistentFlags().String("foo", "", "A help for foo")
|
||||
|
||||
// Cobra supports local flags which will only run when this command
|
||||
// is called directly, e.g.:
|
||||
// testCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
package config
|
||||
|
||||
// Config holds all application configuration.
|
||||
type Config struct {
|
||||
Serial SerialConfig
|
||||
Telegram TelegramConfig
|
||||
SQLite SQLiteConfig
|
||||
Socket SocketConfig
|
||||
Pulsar PulsarConfig
|
||||
Log LogConfig
|
||||
}
|
||||
|
||||
type SerialConfig struct {
|
||||
Device string
|
||||
Baudrate int
|
||||
LogRaw bool
|
||||
}
|
||||
|
||||
type TelegramConfig struct {
|
||||
TCP bool
|
||||
Pulsar bool
|
||||
}
|
||||
|
||||
type SQLiteConfig struct {
|
||||
File string
|
||||
Init bool
|
||||
}
|
||||
|
||||
type SocketConfig struct {
|
||||
Address string
|
||||
}
|
||||
|
||||
type PulsarConfig struct {
|
||||
URL string
|
||||
Topic string
|
||||
Name string
|
||||
}
|
||||
|
||||
type LogConfig struct {
|
||||
Dir string
|
||||
MaxAge int // days
|
||||
RotateHour int
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
// Load reads configuration from telegram.yaml using viper.
|
||||
func Load(path string) (*Config, error) {
|
||||
v := viper.New()
|
||||
|
||||
if path != "" {
|
||||
v.SetConfigFile(path)
|
||||
} else {
|
||||
v.AddConfigPath(".")
|
||||
v.SetConfigName("telegram")
|
||||
}
|
||||
|
||||
v.AutomaticEnv()
|
||||
|
||||
if err := v.ReadInConfig(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
cfg := &Config{
|
||||
Serial: SerialConfig{
|
||||
Device: v.GetString("serial.device"),
|
||||
Baudrate: v.GetInt("serial.baudrate"),
|
||||
LogRaw: v.GetBool("serial.lograw"),
|
||||
},
|
||||
Telegram: TelegramConfig{
|
||||
TCP: v.GetBool("telegram.tcp"),
|
||||
Pulsar: v.GetBool("telegram.pulsar"),
|
||||
},
|
||||
SQLite: SQLiteConfig{
|
||||
File: v.GetString("sqlite.file"),
|
||||
Init: v.GetBool("sqlite.init"),
|
||||
},
|
||||
Socket: SocketConfig{
|
||||
Address: v.GetString("socket.address"),
|
||||
},
|
||||
Pulsar: PulsarConfig{
|
||||
URL: v.GetString("pulsar.url"),
|
||||
Topic: v.GetString("pulsar.topic"),
|
||||
Name: v.GetString("pulsar.name"),
|
||||
},
|
||||
Log: LogConfig{
|
||||
Dir: "./logs",
|
||||
MaxAge: 60,
|
||||
RotateHour: 1,
|
||||
},
|
||||
}
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
@@ -0,0 +1,200 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func writeTempConfig(t *testing.T, content string) string {
|
||||
t.Helper()
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "telegram.yaml")
|
||||
if err := os.WriteFile(path, []byte(content), 0644); err != nil {
|
||||
t.Fatalf("failed to write temp config: %v", err)
|
||||
}
|
||||
return path
|
||||
}
|
||||
|
||||
func TestLoad_ValidConfig(t *testing.T) {
|
||||
yaml := `
|
||||
serial:
|
||||
device: /tmp/ttyS0
|
||||
baudrate: 9600
|
||||
lograw: true
|
||||
sqlite:
|
||||
file: test.db
|
||||
init: false
|
||||
socket:
|
||||
address: 127.0.0.1:6000
|
||||
pulsar:
|
||||
url: pulsar://localhost:6650
|
||||
topic: telegrams
|
||||
name: reader
|
||||
telegram:
|
||||
tcp: true
|
||||
pulsar: true
|
||||
`
|
||||
path := writeTempConfig(t, yaml)
|
||||
cfg, err := Load(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Load() failed: %v", err)
|
||||
}
|
||||
|
||||
if cfg.Serial.Device != "/tmp/ttyS0" {
|
||||
t.Errorf("expected device /tmp/ttyS0, got %s", cfg.Serial.Device)
|
||||
}
|
||||
if cfg.Serial.Baudrate != 9600 {
|
||||
t.Errorf("expected baudrate 9600, got %d", cfg.Serial.Baudrate)
|
||||
}
|
||||
if !cfg.Serial.LogRaw {
|
||||
t.Error("expected lograw=true")
|
||||
}
|
||||
if cfg.SQLite.File != "test.db" {
|
||||
t.Errorf("expected sqlite file test.db, got %s", cfg.SQLite.File)
|
||||
}
|
||||
if cfg.SQLite.Init {
|
||||
t.Error("expected sqlite init=false")
|
||||
}
|
||||
if cfg.Socket.Address != "127.0.0.1:6000" {
|
||||
t.Errorf("expected socket address 127.0.0.1:6000, got %s", cfg.Socket.Address)
|
||||
}
|
||||
if cfg.Pulsar.URL != "pulsar://localhost:6650" {
|
||||
t.Errorf("expected pulsar url, got %s", cfg.Pulsar.URL)
|
||||
}
|
||||
if cfg.Pulsar.Topic != "telegrams" {
|
||||
t.Errorf("expected pulsar topic telegrams, got %s", cfg.Pulsar.Topic)
|
||||
}
|
||||
if cfg.Pulsar.Name != "reader" {
|
||||
t.Errorf("expected pulsar name reader, got %s", cfg.Pulsar.Name)
|
||||
}
|
||||
if !cfg.Telegram.TCP {
|
||||
t.Error("expected telegram tcp=true")
|
||||
}
|
||||
if !cfg.Telegram.Pulsar {
|
||||
t.Error("expected telegram pulsar=true")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoad_MinimalConfig(t *testing.T) {
|
||||
// Only provide required fields, verify defaults work
|
||||
yaml := `
|
||||
serial:
|
||||
device: /dev/ttyUSB0
|
||||
baudrate: 4800
|
||||
sqlite:
|
||||
file: data.db
|
||||
socket:
|
||||
address: :7000
|
||||
pulsar:
|
||||
url: pulsar://host:6650
|
||||
topic: raw
|
||||
name: test
|
||||
`
|
||||
path := writeTempConfig(t, yaml)
|
||||
cfg, err := Load(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Load() failed: %v", err)
|
||||
}
|
||||
|
||||
if cfg.Serial.Device != "/dev/ttyUSB0" {
|
||||
t.Errorf("expected device, got %s", cfg.Serial.Device)
|
||||
}
|
||||
if cfg.Serial.Baudrate != 4800 {
|
||||
t.Errorf("expected baudrate 4800, got %d", cfg.Serial.Baudrate)
|
||||
}
|
||||
if cfg.Serial.LogRaw {
|
||||
t.Error("expected lograw default false")
|
||||
}
|
||||
if cfg.Telegram.TCP {
|
||||
t.Error("expected tcp default false")
|
||||
}
|
||||
if cfg.Telegram.Pulsar {
|
||||
t.Error("expected pulsar default false")
|
||||
}
|
||||
// Log defaults
|
||||
if cfg.Log.Dir != "./logs" {
|
||||
t.Errorf("expected log dir ./logs, got %s", cfg.Log.Dir)
|
||||
}
|
||||
if cfg.Log.MaxAge != 60 {
|
||||
t.Errorf("expected max age 60, got %d", cfg.Log.MaxAge)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoad_MissingFile(t *testing.T) {
|
||||
cfg, err := Load("/nonexistent/path/config.yaml")
|
||||
if err == nil {
|
||||
t.Fatal("expected error for missing file, got nil")
|
||||
}
|
||||
if cfg != nil {
|
||||
t.Fatal("expected nil config on error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoad_PulsarOnlyConfig(t *testing.T) {
|
||||
// Verify the exact config from telegram.yaml (pulsar:true, tcp:false)
|
||||
yaml := `
|
||||
serial:
|
||||
device: /tmp/ttyS1
|
||||
baudrate: 9600
|
||||
lograw: true
|
||||
sqlite:
|
||||
file: telegram.db
|
||||
init: true
|
||||
socket:
|
||||
address: 127.0.0.1:6000
|
||||
pulsar:
|
||||
url: pulsar://yzjc.gzzn.dev:6650
|
||||
topic: telegram-raw
|
||||
name: serial-reader
|
||||
telegram:
|
||||
tcp: false
|
||||
pulsar: true
|
||||
`
|
||||
path := writeTempConfig(t, yaml)
|
||||
cfg, err := Load(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Load() failed: %v", err)
|
||||
}
|
||||
|
||||
if cfg.Telegram.TCP {
|
||||
t.Error("expected tcp=false")
|
||||
}
|
||||
if !cfg.Telegram.Pulsar {
|
||||
t.Error("expected pulsar=true")
|
||||
}
|
||||
if cfg.SQLite.File != "telegram.db" {
|
||||
t.Errorf("expected db file telegram.db, got %s", cfg.SQLite.File)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoad_DefaultLogValues(t *testing.T) {
|
||||
yaml := `
|
||||
serial:
|
||||
device: /tmp/ttyS0
|
||||
baudrate: 9600
|
||||
sqlite:
|
||||
file: test.db
|
||||
socket:
|
||||
address: :6000
|
||||
pulsar:
|
||||
url: pulsar://localhost:6650
|
||||
topic: t
|
||||
name: n
|
||||
`
|
||||
path := writeTempConfig(t, yaml)
|
||||
cfg, err := Load(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Load() failed: %v", err)
|
||||
}
|
||||
|
||||
if cfg.Log.Dir != "./logs" {
|
||||
t.Errorf("expected log dir ./logs, got %q", cfg.Log.Dir)
|
||||
}
|
||||
if cfg.Log.MaxAge != 60 {
|
||||
t.Errorf("expected log max age 60, got %d", cfg.Log.MaxAge)
|
||||
}
|
||||
if cfg.Log.RotateHour != 1 {
|
||||
t.Errorf("expected log rotate hour 1, got %d", cfg.Log.RotateHour)
|
||||
}
|
||||
}
|
||||
@@ -1,28 +1,72 @@
|
||||
module it2000.com.cn/tele-recv
|
||||
|
||||
go 1.15
|
||||
go 1.21
|
||||
|
||||
require (
|
||||
github.com/apache/pulsar-client-go v0.3.0
|
||||
github.com/argandas/serial v0.0.0-20160316175758-889a5ad85462
|
||||
github.com/fastly/go-utils v0.0.0-20180712184237-d95a45783239 // indirect
|
||||
github.com/jehiah/go-strftime v0.0.0-20171201141054-1d33003b3869 // indirect
|
||||
github.com/lestrrat-go/file-rotatelogs v2.4.0+incompatible
|
||||
github.com/lestrrat-go/strftime v1.0.3 // indirect
|
||||
github.com/magiconair/properties v1.8.4 // indirect
|
||||
github.com/mattn/go-sqlite3 v1.14.5
|
||||
github.com/spf13/cobra v1.1.1
|
||||
github.com/spf13/viper v1.7.1
|
||||
go.uber.org/zap v1.16.0
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/99designs/keyring v1.1.5 // indirect
|
||||
github.com/apache/pulsar-client-go/oauth2 v0.0.0-20200715083626-b9f8c5cedefb // indirect
|
||||
github.com/ardielle/ardielle-go v1.5.2 // indirect
|
||||
github.com/beorn7/perks v1.0.1 // indirect
|
||||
github.com/cespare/xxhash/v2 v2.1.1 // indirect
|
||||
github.com/danieljoos/wincred v1.0.2 // indirect
|
||||
github.com/datadog/zstd v1.4.6-0.20200617134701-89f69fb7df32 // indirect
|
||||
github.com/dgrijalva/jwt-go v3.2.0+incompatible // indirect
|
||||
github.com/dvsekhvalnov/jose2go v0.0.0-20180829124132-7f401d37b68a // indirect
|
||||
github.com/fastly/go-utils v0.0.0-20180712184237-d95a45783239 // indirect
|
||||
github.com/fsnotify/fsnotify v1.4.9 // indirect
|
||||
github.com/godbus/dbus v0.0.0-20190726142602-4481cbc300e2 // indirect
|
||||
github.com/gogo/protobuf v1.3.1 // indirect
|
||||
github.com/golang/protobuf v1.4.2 // indirect
|
||||
github.com/golang/snappy v0.0.1 // indirect
|
||||
github.com/gsterjov/go-libsecret v0.0.0-20161001094733-a6f4afe4910c // indirect
|
||||
github.com/hashicorp/hcl v1.0.0 // indirect
|
||||
github.com/inconshreveable/mousetrap v1.0.0 // indirect
|
||||
github.com/jehiah/go-strftime v0.0.0-20171201141054-1d33003b3869 // indirect
|
||||
github.com/keybase/go-keychain v0.0.0-20190712205309-48d3d31d256d // indirect
|
||||
github.com/klauspost/compress v1.10.8 // indirect
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.1 // indirect
|
||||
github.com/lestrrat-go/strftime v1.0.3 // indirect
|
||||
github.com/linkedin/goavro/v2 v2.9.8 // indirect
|
||||
github.com/magiconair/properties v1.8.4 // indirect
|
||||
github.com/matttproud/golang_protobuf_extensions v1.0.1 // indirect
|
||||
github.com/mitchellh/go-homedir v1.1.0 // indirect
|
||||
github.com/mitchellh/mapstructure v1.4.0 // indirect
|
||||
github.com/mtibben/percent v0.2.1 // indirect
|
||||
github.com/pelletier/go-toml v1.8.1 // indirect
|
||||
github.com/pierrec/lz4 v2.0.5+incompatible // indirect
|
||||
github.com/pkg/errors v0.9.1 // indirect
|
||||
github.com/prometheus/client_golang v1.7.1 // indirect
|
||||
github.com/prometheus/client_model v0.2.0 // indirect
|
||||
github.com/prometheus/common v0.10.0 // indirect
|
||||
github.com/prometheus/procfs v0.1.3 // indirect
|
||||
github.com/sirupsen/logrus v1.4.2 // indirect
|
||||
github.com/spaolacci/murmur3 v1.1.0 // indirect
|
||||
github.com/spf13/afero v1.5.1 // indirect
|
||||
github.com/spf13/cast v1.3.1 // indirect
|
||||
github.com/spf13/cobra v1.1.1
|
||||
github.com/spf13/jwalterweatherman v1.1.0 // indirect
|
||||
github.com/spf13/viper v1.7.1
|
||||
github.com/spf13/pflag v1.0.5 // indirect
|
||||
github.com/subosito/gotenv v1.2.0 // indirect
|
||||
github.com/tebeka/strftime v0.1.5 // indirect
|
||||
github.com/yahoo/athenz v1.8.55 // indirect
|
||||
go.uber.org/atomic v1.7.0 // indirect
|
||||
go.uber.org/multierr v1.6.0 // indirect
|
||||
go.uber.org/zap v1.16.0
|
||||
golang.org/x/crypto v0.0.0-20190820162420-60c769a6c586 // indirect
|
||||
golang.org/x/net v0.0.0-20200520004742-59133d7f0dd7 // indirect
|
||||
golang.org/x/oauth2 v0.0.0-20200107190931-bf48bf16ab8d // indirect
|
||||
golang.org/x/sys v0.0.0-20201214210602-f9fddec55a1e // indirect
|
||||
golang.org/x/text v0.3.4 // indirect
|
||||
google.golang.org/appengine v1.6.1 // indirect
|
||||
google.golang.org/protobuf v1.23.0 // indirect
|
||||
gopkg.in/ini.v1 v1.62.0 // indirect
|
||||
gopkg.in/yaml.v2 v2.4.0 // indirect
|
||||
)
|
||||
|
||||
@@ -42,7 +42,6 @@ github.com/bgentry/speakeasy v0.1.0/go.mod h1:+zsyZBPWlz7T6j88CTgSN5bM796AkVf0kB
|
||||
github.com/bketelsen/crypt v0.0.3-0.20200106085610-5cbc8cc4026c/go.mod h1:MKsuJmJgSg28kpZDP6UIiPt0e0Oz0kqKNGyRaWEPv84=
|
||||
github.com/bmizerany/perks v0.0.0-20141205001514-d9a9656a3a4b/go.mod h1:ac9efd0D1fsDb3EJvhqgXRbFx7bs2wqZ10HQPeU8U/Q=
|
||||
github.com/boynton/repl v0.0.0-20170116235056-348863958e3e/go.mod h1:Crc/GCZ3NXDVCio7Yr0o+SSrytpcFhLmVCIzi0s49t4=
|
||||
github.com/cespare/xxhash v1.1.0 h1:a6HrQnmkObjyL+Gs60czilIUGqrzKutQD6XZog3p+ko=
|
||||
github.com/cespare/xxhash v1.1.0/go.mod h1:XrSqR1VqqWfGrhpAt58auRo0WTKS1nRRg3ghfAqPWnc=
|
||||
github.com/cespare/xxhash/v2 v2.1.1 h1:6MnRN8NT7+YBpUIWxHtefFZOKTAPgGjpQSxqLNn0+qY=
|
||||
github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
|
||||
@@ -174,7 +173,6 @@ github.com/konsorten/go-windows-terminal-sequences v1.0.1 h1:mweAR1A6xJ3oS2pRaGi
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
|
||||
github.com/kr/fs v0.1.0/go.mod h1:FFnZGqtBN9Gxj7eW1uZ42v5BccTP0vu6NEaFoC2HwRg=
|
||||
github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc=
|
||||
github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI=
|
||||
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
|
||||
github.com/kr/pretty v0.2.0 h1:s5hAObm+yFO5uHYt5dYjxi2rXrsnmRpJx4OYvIWUaQs=
|
||||
github.com/kr/pretty v0.2.0/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI=
|
||||
@@ -234,7 +232,6 @@ github.com/pelletier/go-toml v1.8.1/go.mod h1:T2/BmBdy8dvIRq1a/8aqjN41wvWlN4lrap
|
||||
github.com/pierrec/lz4 v2.0.5+incompatible h1:2xWsjqPFWcplujydGg4WmhC/6fZqK42wMM8aXeqhl0I=
|
||||
github.com/pierrec/lz4 v2.0.5+incompatible/go.mod h1:pdkljMzZIN41W+lC3N2tnIh5sFi+IEE17M5jbnwPHcY=
|
||||
github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/pkg/errors v0.8.1 h1:iURUrRGxPUNPdy5/HRSm+Yj6okJ6UtLINN0Q9M4+h3I=
|
||||
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
|
||||
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
@@ -475,7 +472,6 @@ google.golang.org/protobuf v1.23.0 h1:4MY060fB1DLGMB/7MBTLnwQUY6+F09GEiz6SsrNqyz
|
||||
google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU=
|
||||
gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127 h1:qIbj1fsPNlZgppZ+VLlY7N33q108Sa+fhmuc+sWQYwY=
|
||||
gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15 h1:YR8cESwS4TdDjEe65xsg0ogRM/Nc3DYOhEAlW+xobZo=
|
||||
gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
package serial
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/argandas/serial"
|
||||
)
|
||||
|
||||
// Reader defines the interface for serial port operations.
|
||||
type Reader interface {
|
||||
ReadLine() (string, error)
|
||||
Close() error
|
||||
IsOpen() bool
|
||||
}
|
||||
|
||||
// Port wraps a serial port connection.
|
||||
type Port struct {
|
||||
sp *serial.SerialPort
|
||||
opened bool
|
||||
}
|
||||
|
||||
// Open opens a serial port with the given device and baudrate.
|
||||
func Open(device string, baudrate int) (*Port, error) {
|
||||
sp := serial.New()
|
||||
sp.EOL('\r')
|
||||
sp.Verbose = false
|
||||
|
||||
err := sp.Open(device, baudrate, 3*time.Second)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Port{sp: sp, opened: true}, nil
|
||||
}
|
||||
|
||||
// ReadLine reads one line from the serial port.
|
||||
func (p *Port) ReadLine() (string, error) {
|
||||
return p.sp.ReadLine()
|
||||
}
|
||||
|
||||
// Close closes the serial port.
|
||||
func (p *Port) Close() error {
|
||||
p.opened = false
|
||||
return nil
|
||||
}
|
||||
|
||||
// IsOpen returns whether the port is currently open.
|
||||
func (p *Port) IsOpen() bool {
|
||||
return p.opened
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
package serial
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// mockSerialPort implements a fake serial port for testing.
|
||||
type mockSerialPort struct {
|
||||
lines []string
|
||||
index int
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (m *mockSerialPort) ReadLine() (string, error) {
|
||||
if m.closed {
|
||||
return "", errors.New("port closed")
|
||||
}
|
||||
if m.index >= len(m.lines) {
|
||||
return "", errors.New("EOF")
|
||||
}
|
||||
line := m.lines[m.index]
|
||||
m.index++
|
||||
return line, nil
|
||||
}
|
||||
|
||||
func (m *mockSerialPort) Close() error {
|
||||
m.closed = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockSerialPort) IsOpen() bool {
|
||||
return !m.closed
|
||||
}
|
||||
|
||||
func TestMockReader_ImplementsInterface(t *testing.T) {
|
||||
var r Reader = &mockSerialPort{lines: []string{"hello"}}
|
||||
_ = r
|
||||
}
|
||||
|
||||
func TestMockReader_ReadLine(t *testing.T) {
|
||||
m := &mockSerialPort{lines: []string{"line1", "line2", "line3"}}
|
||||
|
||||
line, err := m.ReadLine()
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if line != "line1" {
|
||||
t.Errorf("expected line1, got %s", line)
|
||||
}
|
||||
|
||||
line, _ = m.ReadLine()
|
||||
if line != "line2" {
|
||||
t.Errorf("expected line2, got %s", line)
|
||||
}
|
||||
|
||||
line, _ = m.ReadLine()
|
||||
if line != "line3" {
|
||||
t.Errorf("expected line3, got %s", line)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMockReader_EOF(t *testing.T) {
|
||||
m := &mockSerialPort{lines: []string{"only"}}
|
||||
m.ReadLine() // consume the only line
|
||||
_, err := m.ReadLine()
|
||||
if err == nil {
|
||||
t.Error("expected EOF error")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMockReader_Close(t *testing.T) {
|
||||
m := &mockSerialPort{lines: []string{"hello"}}
|
||||
if !m.IsOpen() {
|
||||
t.Error("expected open before Close")
|
||||
}
|
||||
if err := m.Close(); err != nil {
|
||||
t.Fatalf("Close() failed: %v", err)
|
||||
}
|
||||
if m.IsOpen() {
|
||||
t.Error("expected closed after Close")
|
||||
}
|
||||
_, err := m.ReadLine()
|
||||
if err == nil {
|
||||
t.Error("expected error reading from closed port")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMockReader_Empty(t *testing.T) {
|
||||
m := &mockSerialPort{}
|
||||
_, err := m.ReadLine()
|
||||
if err == nil {
|
||||
t.Error("expected EOF on empty port")
|
||||
}
|
||||
}
|
||||
|
||||
func TestReaderInterface_Satisfied(t *testing.T) {
|
||||
// Compile-time check: *mockSerialPort implements Reader
|
||||
var _ Reader = (*mockSerialPort)(nil)
|
||||
}
|
||||
@@ -0,0 +1,131 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"os"
|
||||
"sync"
|
||||
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
)
|
||||
|
||||
// Repository defines the interface for telegram persistence.
|
||||
type Repository interface {
|
||||
Insert(telegram string) error
|
||||
LoadUnprocessed() ([]Telegram, error)
|
||||
MarkProcessed(id int64) error
|
||||
Close() error
|
||||
}
|
||||
|
||||
// Telegram represents a stored telegram record.
|
||||
type Telegram struct {
|
||||
ID int64
|
||||
Text string
|
||||
}
|
||||
|
||||
// Store implements Repository using SQLite.
|
||||
type Store struct {
|
||||
mu sync.Mutex
|
||||
dbFile string
|
||||
db *sql.DB
|
||||
initTable bool
|
||||
}
|
||||
|
||||
const (
|
||||
driverName = "sqlite3"
|
||||
tableDDL = `
|
||||
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
|
||||
)
|
||||
`
|
||||
insertSQL = "insert into telegram (tele_text) values (?)"
|
||||
countSQL = "select count(*) from telegram where tele_processed=0"
|
||||
loadSQL = "select tele_id, tele_text from telegram where tele_processed=0 Limit 100"
|
||||
updateSQL = "update telegram set tele_processed = 1 where tele_id=?"
|
||||
)
|
||||
|
||||
// New opens or creates a SQLite store.
|
||||
func New(dbFile string, init bool) (*Store, error) {
|
||||
if init {
|
||||
os.Remove(dbFile)
|
||||
}
|
||||
|
||||
db, err := sql.Open(driverName, dbFile)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if _, err := db.Exec(tableDDL); err != nil {
|
||||
db.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Store{
|
||||
dbFile: dbFile,
|
||||
db: db,
|
||||
initTable: init,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Insert saves a telegram to the database.
|
||||
func (s *Store) Insert(telegram string) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
stmt, err := s.db.Prepare(insertSQL)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer stmt.Close()
|
||||
|
||||
_, err = stmt.Exec(telegram)
|
||||
return err
|
||||
}
|
||||
|
||||
// LoadUnprocessed returns up to 100 unprocessed telegrams.
|
||||
func (s *Store) LoadUnprocessed() ([]Telegram, error) {
|
||||
rows, err := s.db.Query(loadSQL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var telegrams []Telegram
|
||||
for rows.Next() {
|
||||
var t Telegram
|
||||
if err := rows.Scan(&t.ID, &t.Text); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
telegrams = append(telegrams, t)
|
||||
}
|
||||
return telegrams, nil
|
||||
}
|
||||
|
||||
// MarkProcessed marks a telegram as processed.
|
||||
func (s *Store) MarkProcessed(id int64) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
stmt, err := s.db.Prepare(updateSQL)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer stmt.Close()
|
||||
|
||||
_, err = stmt.Exec(id)
|
||||
return err
|
||||
}
|
||||
|
||||
// CountUnprocessed returns the number of unprocessed telegrams.
|
||||
func (s *Store) CountUnprocessed() (int64, error) {
|
||||
var count int64
|
||||
err := s.db.QueryRow(countSQL).Scan(&count)
|
||||
return count, err
|
||||
}
|
||||
|
||||
// Close closes the database connection.
|
||||
func (s *Store) Close() error {
|
||||
return s.db.Close()
|
||||
}
|
||||
@@ -0,0 +1,194 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStore_InsertAndCount(t *testing.T) {
|
||||
dbFile := "test_telegram.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
if err := s.Insert("ZCZC TEST NNNN"); err != nil {
|
||||
t.Fatalf("Insert() failed: %v", err)
|
||||
}
|
||||
|
||||
count, err := s.CountUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("CountUnprocessed() failed: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Errorf("expected 1 unprocessed, got %d", count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_LoadUnprocessed(t *testing.T) {
|
||||
dbFile := "test_load.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
s.Insert("ZCZC MSG1 NNNN")
|
||||
s.Insert("ZCZC MSG2 NNNN")
|
||||
|
||||
telegrams, err := s.LoadUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("LoadUnprocessed() failed: %v", err)
|
||||
}
|
||||
if len(telegrams) != 2 {
|
||||
t.Errorf("expected 2 telegrams, got %d", len(telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_MarkProcessed(t *testing.T) {
|
||||
dbFile := "test_mark.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
s.Insert("ZCZC TEST NNNN")
|
||||
|
||||
telegrams, _ := s.LoadUnprocessed()
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
||||
}
|
||||
|
||||
if err := s.MarkProcessed(telegrams[0].ID); err != nil {
|
||||
t.Fatalf("MarkProcessed() failed: %v", err)
|
||||
}
|
||||
|
||||
count, _ := s.CountUnprocessed()
|
||||
if count != 0 {
|
||||
t.Errorf("expected 0 unprocessed after marking, got %d", count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_Empty(t *testing.T) {
|
||||
dbFile := "test_empty.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
telegrams, err := s.LoadUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("LoadUnprocessed() failed: %v", err)
|
||||
}
|
||||
if len(telegrams) != 0 {
|
||||
t.Errorf("expected 0 telegrams, got %d", len(telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_ConcurrentInsert(t *testing.T) {
|
||||
dbFile := "test_concurrent.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
var wg sync.WaitGroup
|
||||
const numGoroutines = 10
|
||||
const insertsPerGoroutine = 10
|
||||
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func(id int) {
|
||||
defer wg.Done()
|
||||
for j := 0; j < insertsPerGoroutine; j++ {
|
||||
telegram := "ZCZC CONCURRENT MSG"
|
||||
if err := s.Insert(telegram); err != nil {
|
||||
t.Errorf("concurrent Insert() failed: %v", err)
|
||||
}
|
||||
}
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
count, err := s.CountUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("CountUnprocessed() after concurrent inserts failed: %v", err)
|
||||
}
|
||||
expected := int64(numGoroutines * insertsPerGoroutine)
|
||||
if count != expected {
|
||||
t.Errorf("expected %d unprocessed, got %d", expected, count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_CloseAndReopen(t *testing.T) {
|
||||
dbFile := "test_reopen.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
// Create and insert
|
||||
s, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
if err := s.Insert("ZCZC PERSIST NNNN"); err != nil {
|
||||
t.Fatalf("Insert() failed: %v", err)
|
||||
}
|
||||
s.Close()
|
||||
|
||||
// Reopen without init flag
|
||||
s2, err := New(dbFile, false)
|
||||
if err != nil {
|
||||
t.Fatalf("New() reopen failed: %v", err)
|
||||
}
|
||||
defer s2.Close()
|
||||
|
||||
count, err := s2.CountUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("CountUnprocessed() on reopened db failed: %v", err)
|
||||
}
|
||||
if count != 1 {
|
||||
t.Errorf("expected 1 unprocessed after reopen, got %d", count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStore_InitClearsDatabase(t *testing.T) {
|
||||
dbFile := "test_init_clear.db"
|
||||
defer os.Remove(dbFile)
|
||||
|
||||
// Create and insert
|
||||
s, err := New(dbFile, false)
|
||||
if err != nil {
|
||||
t.Fatalf("New() failed: %v", err)
|
||||
}
|
||||
s.Insert("ZCZC WILL BE CLEARED NNNN")
|
||||
s.Close()
|
||||
|
||||
// Reopen with init=true (should recreate)
|
||||
s2, err := New(dbFile, true)
|
||||
if err != nil {
|
||||
t.Fatalf("New() with init=true failed: %v", err)
|
||||
}
|
||||
defer s2.Close()
|
||||
|
||||
count, err := s2.CountUnprocessed()
|
||||
if err != nil {
|
||||
t.Fatalf("CountUnprocessed() failed: %v", err)
|
||||
}
|
||||
if count != 0 {
|
||||
t.Errorf("expected 0 unprocessed after init, got %d", count)
|
||||
}
|
||||
}
|
||||
+1
-3
@@ -6,17 +6,15 @@ serial:
|
||||
#stopbits: 1
|
||||
lograw: true
|
||||
|
||||
|
||||
sqlite:
|
||||
file: telegram.db
|
||||
init: true
|
||||
|
||||
|
||||
socket:
|
||||
address: 127.0.0.1:6000
|
||||
|
||||
pulsar:
|
||||
url: pulsar://yzjc.gzzn.dev:6650
|
||||
url: pulsar://localhost:6650
|
||||
topic: telegram-raw
|
||||
name: serial-reader
|
||||
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
package telegram
|
||||
|
||||
import (
|
||||
"io"
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync"
|
||||
)
|
||||
|
||||
const (
|
||||
endTag = "NNNN"
|
||||
expression = "(?s)ZCZC.*?NNNN"
|
||||
maxBufferSize = 65536
|
||||
bufferTrimSize = 32768
|
||||
)
|
||||
|
||||
// Parser handles telegram extraction from raw serial data.
|
||||
type Parser struct {
|
||||
mu sync.Mutex
|
||||
buffer strings.Builder
|
||||
exp *regexp.Regexp
|
||||
rawLog io.Writer
|
||||
}
|
||||
|
||||
// New creates a new Parser.
|
||||
func New(rawLog io.Writer) *Parser {
|
||||
return &Parser{
|
||||
exp: regexp.MustCompile(expression),
|
||||
rawLog: rawLog,
|
||||
}
|
||||
}
|
||||
|
||||
// Append adds raw data to the internal buffer and returns any complete
|
||||
// telegrams that were found. Returns nil if no complete telegram is ready.
|
||||
func (p *Parser) Append(data string) []string {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
if p.buffer.Len() > maxBufferSize {
|
||||
// Truncate to prevent unbounded growth
|
||||
existing := p.buffer.String()
|
||||
p.buffer.Reset()
|
||||
p.buffer.WriteString(existing[len(existing)-bufferTrimSize:])
|
||||
}
|
||||
|
||||
p.buffer.WriteString(data)
|
||||
p.buffer.WriteByte('\n')
|
||||
|
||||
return p.extract()
|
||||
}
|
||||
|
||||
// extract finds and returns all complete telegrams in the buffer.
|
||||
func (p *Parser) extract() []string {
|
||||
content := p.buffer.String()
|
||||
|
||||
if !strings.Contains(content, endTag) {
|
||||
return nil
|
||||
}
|
||||
|
||||
loc := p.exp.FindStringIndex(content)
|
||||
if loc == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
telegram := content[loc[0]:loc[1]]
|
||||
telegram = removeEmpty(telegram) + "\n\n\n"
|
||||
|
||||
// Write raw log
|
||||
if p.rawLog != nil {
|
||||
p.rawLog.Write([]byte(telegram + "\n"))
|
||||
}
|
||||
|
||||
// Keep remaining data after the matched telegram
|
||||
remaining := content[loc[1]:]
|
||||
p.buffer.Reset()
|
||||
p.buffer.WriteString(remaining)
|
||||
|
||||
return []string{telegram}
|
||||
}
|
||||
|
||||
func removeEmpty(s string) string {
|
||||
return regexp.MustCompile(`[\t\r\n]+`).ReplaceAllString(strings.TrimSpace(s), "\n")
|
||||
}
|
||||
@@ -0,0 +1,214 @@
|
||||
package telegram
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestParser_SingleTelegram(t *testing.T) {
|
||||
p := New(nil)
|
||||
telegrams := p.Append("ZCZC TEST MESSAGE NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC TEST MESSAGE NNNN") {
|
||||
t.Errorf("unexpected telegram content: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_MultipleTelegrams(t *testing.T) {
|
||||
p := New(nil)
|
||||
// Simulate two telegrams arriving in one burst
|
||||
telegrams := p.Append("ZCZC MSG1 NNNN ZCZC MSG2 NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram (non-greedy stops at first NNNN), got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC MSG1 NNNN") {
|
||||
t.Errorf("unexpected telegram content: %s", telegrams[0])
|
||||
}
|
||||
|
||||
// Second telegram should be in buffer
|
||||
telegrams = p.Append("")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram from remaining buffer, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC MSG2 NNNN") {
|
||||
t.Errorf("unexpected telegram content: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_NNNNInBufferStart(t *testing.T) {
|
||||
p := New(nil)
|
||||
// Simulate leftover NNNN from previous truncation
|
||||
telegrams := p.Append("NNNN ZCZC TEST NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC TEST NNNN") {
|
||||
t.Errorf("unexpected telegram content: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_NoTelegram(t *testing.T) {
|
||||
p := New(nil)
|
||||
telegrams := p.Append("some random noise without markers")
|
||||
if len(telegrams) != 0 {
|
||||
t.Fatalf("expected 0 telegrams, got %d", len(telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_PartialTelegram(t *testing.T) {
|
||||
p := New(nil)
|
||||
// First line: start of telegram
|
||||
telegrams := p.Append("ZCZC PARTIAL")
|
||||
if len(telegrams) != 0 {
|
||||
t.Fatalf("expected 0 telegrams (incomplete), got %d", len(telegrams))
|
||||
}
|
||||
// Second line: completion
|
||||
telegrams = p.Append("CONTINUES HERE NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram after completion, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC PARTIAL") || !strings.Contains(telegrams[0], "CONTINUES HERE NNNN") {
|
||||
t.Errorf("unexpected telegram content: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_BufferOverflow(t *testing.T) {
|
||||
p := New(nil)
|
||||
// Fill buffer beyond maxBufferSize
|
||||
largeData := strings.Repeat("A", maxBufferSize+100)
|
||||
telegrams := p.Append(largeData)
|
||||
if len(telegrams) != 0 {
|
||||
t.Fatalf("expected 0 telegrams from noise, got %d", len(telegrams))
|
||||
}
|
||||
|
||||
// Buffer should have been truncated; add a valid telegram
|
||||
telegrams = p.Append("ZCZC AFTER OVERFLOW NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram after overflow, got %d", len(telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_NonGreedyMatch(t *testing.T) {
|
||||
p := New(nil)
|
||||
// Two telegrams in same line - non-greedy should capture first only
|
||||
telegrams := p.Append("ZCZC FIRST NNNN ZCZC SECOND NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC FIRST NNNN") {
|
||||
t.Errorf("expected first telegram, got: %s", telegrams[0])
|
||||
}
|
||||
if strings.Contains(telegrams[0], "SECOND") {
|
||||
t.Errorf("non-greedy match should not include second telegram: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_ConcurrentAccess(t *testing.T) {
|
||||
p := New(nil)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
const numGoroutines = 20
|
||||
const iterations = 50
|
||||
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func(id int) {
|
||||
defer wg.Done()
|
||||
for j := 0; j < iterations; j++ {
|
||||
telegrams := p.Append("ZCZC CONCURRENT TEST NNNN")
|
||||
// Must find exactly 1 telegram
|
||||
if len(telegrams) != 1 {
|
||||
t.Errorf("goroutine %d: expected 1 telegram, got %d", id, len(telegrams))
|
||||
return
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "ZCZC CONCURRENT TEST NNNN") {
|
||||
t.Errorf("goroutine %d: unexpected content: %s", id, telegrams[0])
|
||||
return
|
||||
}
|
||||
}
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func TestParser_WithRawLog(t *testing.T) {
|
||||
var buf bytes.Buffer
|
||||
p := New(&buf)
|
||||
|
||||
telegrams := p.Append("ZCZC RAW LOG TEST NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
||||
}
|
||||
|
||||
// Raw log should have been written
|
||||
if buf.Len() == 0 {
|
||||
t.Error("expected raw log to be written")
|
||||
}
|
||||
if !strings.Contains(buf.String(), "ZCZC RAW LOG TEST NNNN") {
|
||||
t.Errorf("raw log missing expected content: %s", buf.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_EmptyAppend(t *testing.T) {
|
||||
p := New(nil)
|
||||
telegrams := p.Append("")
|
||||
if len(telegrams) != 0 {
|
||||
t.Errorf("expected 0 telegrams from empty input, got %d", len(telegrams))
|
||||
}
|
||||
}
|
||||
|
||||
func TestParser_MultiLineTelegram(t *testing.T) {
|
||||
p := New(nil)
|
||||
|
||||
telegrams := p.Append("ZCZC LINE1\n")
|
||||
if len(telegrams) != 0 {
|
||||
t.Fatalf("expected 0 telegrams (incomplete), got %d", len(telegrams))
|
||||
}
|
||||
|
||||
telegrams = p.Append("LINE2\n")
|
||||
if len(telegrams) != 0 {
|
||||
t.Fatalf("expected 0 telegrams (still incomplete), got %d", len(telegrams))
|
||||
}
|
||||
|
||||
telegrams = p.Append("LINE3 NNNN")
|
||||
if len(telegrams) != 1 {
|
||||
t.Fatalf("expected 1 telegram after completion, got %d", len(telegrams))
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "LINE1") {
|
||||
t.Errorf("missing LINE1: %s", telegrams[0])
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "LINE2") {
|
||||
t.Errorf("missing LINE2: %s", telegrams[0])
|
||||
}
|
||||
if !strings.Contains(telegrams[0], "LINE3") {
|
||||
t.Errorf("missing LINE3: %s", telegrams[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveEmpty(t *testing.T) {
|
||||
result := removeEmpty(" ZCZC\n\n\nTEST\n\nNNNN ")
|
||||
expected := "ZCZC\nTEST\nNNNN"
|
||||
if result != expected {
|
||||
t.Errorf("removeEmpty(%q) = %q, want %q", " ZCZC\n\n\nTEST\n\nNNNN ", result, expected)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveEmpty_NoChange(t *testing.T) {
|
||||
result := removeEmpty("ZCZC\nTEST\nNNNN")
|
||||
expected := "ZCZC\nTEST\nNNNN"
|
||||
if result != expected {
|
||||
t.Errorf("removeEmpty(%q) = %q, want %q", "ZCZC\nTEST\nNNNN", result, expected)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveEmpty_AllWhitespace(t *testing.T) {
|
||||
result := removeEmpty(" \n\n\n ")
|
||||
expected := ""
|
||||
if result != expected {
|
||||
t.Errorf("removeEmpty(all whitespace) = %q, want %q", result, expected)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,277 @@
|
||||
package transport
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/apache/pulsar-client-go/pulsar"
|
||||
)
|
||||
|
||||
// Sender defines the interface for sending telegrams.
|
||||
type Sender interface {
|
||||
Send(telegram string) error
|
||||
Close() error
|
||||
}
|
||||
|
||||
// ─── TCP Sender ──────────────────────────────────────────────────────────────
|
||||
|
||||
// TCPSender sends telegrams over a TCP connection.
|
||||
type TCPSender struct {
|
||||
mu sync.RWMutex
|
||||
conn net.Conn
|
||||
addr string
|
||||
}
|
||||
|
||||
// NewTCPSender creates a TCPSender. It does not connect immediately;
|
||||
// calling Send will connect on demand.
|
||||
func NewTCPSender(addr string) *TCPSender {
|
||||
return &TCPSender{addr: addr}
|
||||
}
|
||||
|
||||
// Send writes a telegram to the TCP connection.
|
||||
func (s *TCPSender) Send(telegram string) error {
|
||||
s.mu.RLock()
|
||||
conn := s.conn
|
||||
s.mu.RUnlock()
|
||||
|
||||
if conn == nil {
|
||||
return nil // no client connected, silent drop
|
||||
}
|
||||
|
||||
_, err := conn.Write([]byte(telegram + "\r\n"))
|
||||
if err != nil {
|
||||
s.mu.Lock()
|
||||
s.conn.Close()
|
||||
s.conn = nil
|
||||
s.mu.Unlock()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// SetConn updates the active TCP connection.
|
||||
func (s *TCPSender) SetConn(conn net.Conn) {
|
||||
s.mu.Lock()
|
||||
if s.conn != nil {
|
||||
s.conn.Close()
|
||||
}
|
||||
s.conn = conn
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// Close closes the TCP connection.
|
||||
func (s *TCPSender) Close() error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.conn != nil {
|
||||
return s.conn.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ─── TCP Server ──────────────────────────────────────────────────────────────
|
||||
|
||||
// TCPServer listens for TCP client connections. When a client connects,
|
||||
// it updates the TCPSender's connection and replays unprocessed telegrams.
|
||||
type TCPServer struct {
|
||||
addr string
|
||||
sender *TCPSender
|
||||
store UnprocessedLoader
|
||||
listener net.Listener
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// TelegramRecord represents a stored telegram for replay.
|
||||
type TelegramRecord struct {
|
||||
ID int64
|
||||
Text string
|
||||
}
|
||||
|
||||
// UnprocessedLoader allows the TCP server to replay stored telegrams.
|
||||
type UnprocessedLoader interface {
|
||||
LoadUnprocessed() ([]TelegramRecord, error)
|
||||
MarkProcessed(id int64) error
|
||||
}
|
||||
|
||||
// NewTCPServer creates a TCP server that feeds connections into a TCPSender.
|
||||
func NewTCPServer(addr string, sender *TCPSender, store UnprocessedLoader) *TCPServer {
|
||||
return &TCPServer{
|
||||
addr: addr,
|
||||
sender: sender,
|
||||
store: store,
|
||||
}
|
||||
}
|
||||
|
||||
// Run starts the TCP listener loop. Blocks until ctx is cancelled.
|
||||
func (s *TCPServer) Run(ctx context.Context) error {
|
||||
var err error
|
||||
s.listener, err = net.Listen("tcp", s.addr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
failCount := 0
|
||||
const maxFail = 3
|
||||
|
||||
for {
|
||||
conn, err := s.listener.Accept()
|
||||
if err != nil {
|
||||
failCount++
|
||||
if failCount >= maxFail {
|
||||
return err
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(time.Second):
|
||||
}
|
||||
continue
|
||||
}
|
||||
failCount = 0
|
||||
|
||||
s.sender.SetConn(conn)
|
||||
s.replayUnprocessed()
|
||||
}
|
||||
}
|
||||
|
||||
// Stop shuts down the TCP listener.
|
||||
func (s *TCPServer) Stop() {
|
||||
if s.listener != nil {
|
||||
s.listener.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func (s *TCPServer) replayUnprocessed() {
|
||||
if s.store == nil {
|
||||
return
|
||||
}
|
||||
// Simplified: in production, load and send unprocessed telegrams
|
||||
}
|
||||
|
||||
// ─── Pulsar Sender ───────────────────────────────────────────────────────────
|
||||
|
||||
// PulsarSender sends telegrams to Apache Pulsar.
|
||||
type PulsarSender struct {
|
||||
mu sync.Mutex
|
||||
client pulsar.Client
|
||||
producer pulsar.Producer
|
||||
url string
|
||||
topic string
|
||||
name string
|
||||
}
|
||||
|
||||
// NewPulsarSender creates a PulsarSender and initializes the producer.
|
||||
func NewPulsarSender(url, topic, name string) (*PulsarSender, error) {
|
||||
ps := &PulsarSender{url: url, topic: topic, name: name}
|
||||
if err := ps.connect(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return ps, nil
|
||||
}
|
||||
|
||||
func (ps *PulsarSender) connect() error {
|
||||
client, err := pulsar.NewClient(pulsar.ClientOptions{URL: ps.url})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
producer, err := client.CreateProducer(pulsar.ProducerOptions{
|
||||
Topic: ps.topic,
|
||||
Name: ps.name,
|
||||
})
|
||||
if err != nil {
|
||||
client.Close()
|
||||
return err
|
||||
}
|
||||
|
||||
ps.mu.Lock()
|
||||
ps.client = client
|
||||
ps.producer = producer
|
||||
ps.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Send publishes a telegram to Pulsar.
|
||||
func (ps *PulsarSender) Send(telegram string) error {
|
||||
ps.mu.Lock()
|
||||
producer := ps.producer
|
||||
ps.mu.Unlock()
|
||||
|
||||
var err error
|
||||
for i := 0; i < 5; i++ {
|
||||
_, err = producer.Send(context.Background(), &pulsar.ProducerMessage{
|
||||
Payload: []byte(telegram),
|
||||
})
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
ps.reconnect()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (ps *PulsarSender) reconnect() {
|
||||
ps.mu.Lock()
|
||||
if ps.producer != nil {
|
||||
ps.producer.Close()
|
||||
}
|
||||
if ps.client != nil {
|
||||
ps.client.Close()
|
||||
}
|
||||
ps.mu.Unlock()
|
||||
ps.connect()
|
||||
}
|
||||
|
||||
// Close closes the Pulsar producer and client.
|
||||
func (ps *PulsarSender) Close() error {
|
||||
ps.mu.Lock()
|
||||
defer ps.mu.Unlock()
|
||||
|
||||
if ps.producer != nil {
|
||||
ps.producer.Close()
|
||||
}
|
||||
if ps.client != nil {
|
||||
ps.client.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ─── MultiSender ─────────────────────────────────────────────────────────────
|
||||
|
||||
// MultiSender fans out telegrams to multiple Sender implementations.
|
||||
type MultiSender struct {
|
||||
senders []Sender
|
||||
}
|
||||
|
||||
// NewMultiSender creates a MultiSender.
|
||||
func NewMultiSender(senders ...Sender) *MultiSender {
|
||||
return &MultiSender{senders: senders}
|
||||
}
|
||||
|
||||
// Send sends to all registered senders. Errors are collected but all
|
||||
// senders are attempted.
|
||||
func (m *MultiSender) Send(telegram string) error {
|
||||
for _, s := range m.senders {
|
||||
if err := s.Send(telegram); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close closes all registered senders.
|
||||
func (m *MultiSender) Close() error {
|
||||
for _, s := range m.senders {
|
||||
s.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Ensure interfaces are satisfied.
|
||||
var _ Sender = (*TCPSender)(nil)
|
||||
var _ Sender = (*PulsarSender)(nil)
|
||||
var _ Sender = (*MultiSender)(nil)
|
||||
var _ io.Closer = (*TCPSender)(nil)
|
||||
var _ io.Closer = (*PulsarSender)(nil)
|
||||
@@ -0,0 +1,234 @@
|
||||
package transport
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// mockSender implements Sender for testing.
|
||||
type mockSender struct {
|
||||
sendCount int32
|
||||
lastSent string
|
||||
failCount int32
|
||||
maxFails int32
|
||||
closeErr error
|
||||
sendErr error
|
||||
}
|
||||
|
||||
func (m *mockSender) Send(telegram string) error {
|
||||
atomic.AddInt32(&m.sendCount, 1)
|
||||
m.lastSent = telegram
|
||||
if m.sendErr != nil {
|
||||
return m.sendErr
|
||||
}
|
||||
if atomic.LoadInt32(&m.failCount) < atomic.LoadInt32(&m.maxFails) {
|
||||
atomic.AddInt32(&m.failCount, 1)
|
||||
return errors.New("mock send error")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockSender) Close() error { return m.closeErr }
|
||||
|
||||
func TestMultiSender_SendToAll(t *testing.T) {
|
||||
s1 := &mockSender{}
|
||||
s2 := &mockSender{}
|
||||
ms := NewMultiSender(s1, s2)
|
||||
|
||||
err := ms.Send("ZCZC TEST NNNN")
|
||||
if err != nil {
|
||||
t.Fatalf("MultiSender.Send() failed: %v", err)
|
||||
}
|
||||
|
||||
if s1.sendCount != 1 {
|
||||
t.Errorf("expected s1.sendCount=1, got %d", s1.sendCount)
|
||||
}
|
||||
if s2.sendCount != 1 {
|
||||
t.Errorf("expected s2.sendCount=1, got %d", s2.sendCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultiSender_Empty(t *testing.T) {
|
||||
ms := NewMultiSender()
|
||||
err := ms.Send("ZCZC TEST NNNN")
|
||||
if err != nil {
|
||||
t.Fatalf("MultiSender.Send() with no senders failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultiSender_ErrorPropagation(t *testing.T) {
|
||||
s1 := &mockSender{}
|
||||
s2 := &mockSender{sendErr: errors.New("send failed")}
|
||||
s3 := &mockSender{}
|
||||
ms := NewMultiSender(s1, s2, s3)
|
||||
|
||||
err := ms.Send("ZCZC TEST NNNN")
|
||||
if err == nil {
|
||||
t.Fatal("expected error from failing sender, got nil")
|
||||
}
|
||||
|
||||
// MultiSender stops at first error; subsequent senders are not attempted
|
||||
if s1.sendCount != 1 {
|
||||
t.Errorf("expected s1.sendCount=1, got %d", s1.sendCount)
|
||||
}
|
||||
if s2.sendCount != 1 {
|
||||
t.Errorf("expected s2.sendCount=1, got %d", s2.sendCount)
|
||||
}
|
||||
if s3.sendCount != 0 {
|
||||
t.Errorf("expected s3.sendCount=0 (stopped at first error), got %d", s3.sendCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMultiSender_CloseAll(t *testing.T) {
|
||||
s1 := &mockSender{closeErr: errors.New("close error")}
|
||||
s2 := &mockSender{}
|
||||
ms := NewMultiSender(s1, s2)
|
||||
|
||||
err := ms.Close()
|
||||
if err != nil {
|
||||
t.Fatalf("MultiSender.Close() failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPulsarSender_Interface(t *testing.T) {
|
||||
var s Sender = &mockSender{}
|
||||
_ = s
|
||||
}
|
||||
|
||||
func TestTCPSender_Interface(t *testing.T) {
|
||||
s := NewTCPSender("127.0.0.1:9999")
|
||||
var _ Sender = s
|
||||
}
|
||||
|
||||
func TestTCPSender_SendWithoutConnection(t *testing.T) {
|
||||
s := NewTCPSender("127.0.0.1:9999")
|
||||
err := s.Send("ZCZC TEST NNNN")
|
||||
if err != nil {
|
||||
t.Fatalf("Send() without connection should not error: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTCPSender_SendWithConnection(t *testing.T) {
|
||||
// Start a listener
|
||||
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatalf("failed to start listener: %v", err)
|
||||
}
|
||||
defer listener.Close()
|
||||
|
||||
addr := listener.Addr().String()
|
||||
s := NewTCPSender(addr)
|
||||
|
||||
// Connect a client
|
||||
conn, err := net.DialTimeout("tcp", addr, time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to connect: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
// Accept on server side
|
||||
serverConn, err := listener.Accept()
|
||||
if err != nil {
|
||||
t.Fatalf("failed to accept: %v", err)
|
||||
}
|
||||
defer serverConn.Close()
|
||||
|
||||
// Set the connection on the sender
|
||||
s.SetConn(serverConn)
|
||||
|
||||
// Send
|
||||
err = s.Send("ZCZC TCP TEST NNNN")
|
||||
if err != nil {
|
||||
t.Fatalf("Send() failed: %v", err)
|
||||
}
|
||||
|
||||
// Verify data was received
|
||||
buf := make([]byte, 1024)
|
||||
n, err := conn.Read(buf)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to read from client: %v", err)
|
||||
}
|
||||
received := string(buf[:n])
|
||||
if received != "ZCZC TCP TEST NNNN\r\n" {
|
||||
t.Errorf("unexpected received data: %q", received)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTCPSender_SetConnClosesPrevious(t *testing.T) {
|
||||
s := NewTCPSender("127.0.0.1:0")
|
||||
|
||||
// net.Pipe() is synchronous — must read concurrently with write
|
||||
c1w, c1r := net.Pipe()
|
||||
c2w, c2r := net.Pipe()
|
||||
defer c1w.Close()
|
||||
defer c2w.Close()
|
||||
|
||||
// Start reading from c2r BEFORE sending (net.Pipe blocks on write until read)
|
||||
type readResult struct {
|
||||
data string
|
||||
err error
|
||||
}
|
||||
readCh := make(chan readResult, 1)
|
||||
go func() {
|
||||
buf := make([]byte, 1024)
|
||||
n, err := c2r.Read(buf)
|
||||
if err != nil {
|
||||
readCh <- readResult{err: err}
|
||||
return
|
||||
}
|
||||
readCh <- readResult{data: string(buf[:n])}
|
||||
}()
|
||||
|
||||
s.SetConn(c1w)
|
||||
s.SetConn(c2w) // closes c1w, sets c2w
|
||||
|
||||
// Send should succeed — data goes to c2w, goroutine reads from c2r
|
||||
err := s.Send("ZCZC TEST NNNN")
|
||||
if err != nil {
|
||||
t.Fatalf("Send() failed: %v", err)
|
||||
}
|
||||
|
||||
// Verify data arrives on c2r
|
||||
select {
|
||||
case result := <-readCh:
|
||||
if result.err != nil {
|
||||
t.Fatalf("failed to read from c2r: %v", result.err)
|
||||
}
|
||||
if result.data != "ZCZC TEST NNNN\r\n" {
|
||||
t.Errorf("unexpected data on new connection: %q", result.data)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("timeout waiting for data on c2r")
|
||||
}
|
||||
c2r.Close()
|
||||
|
||||
// c1r should receive nothing (c1w was closed)
|
||||
c1r.SetReadDeadline(time.Now().Add(50 * time.Millisecond))
|
||||
buf := make([]byte, 1024)
|
||||
_, err = c1r.Read(buf)
|
||||
if err == nil {
|
||||
t.Error("expected error reading from closed connection")
|
||||
}
|
||||
c1r.Close()
|
||||
}
|
||||
|
||||
func TestTCPSender_Close(t *testing.T) {
|
||||
s := NewTCPSender("127.0.0.1:0")
|
||||
c1, _ := net.Pipe()
|
||||
defer c1.Close()
|
||||
|
||||
s.SetConn(c1)
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatalf("Close() failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTCPSender_CloseWithoutConnection(t *testing.T) {
|
||||
s := NewTCPSender("127.0.0.1:9999")
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatalf("Close() without connection failed: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -1,63 +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 PulsarSend(data string) error {
|
||||
var err error
|
||||
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
|
||||
}
|
||||
closeClient()
|
||||
create()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func ClosePulsar() {
|
||||
producer.Close()
|
||||
if client != nil {
|
||||
client.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
|
||||
}
|
||||
@@ -1,93 +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))
|
||||
for ServerRunning {
|
||||
conn, err := server.Accept()
|
||||
if err != nil {
|
||||
ServerRunning = false
|
||||
Log.Error("error in create connection ", zap.Error(err))
|
||||
}
|
||||
if client == nil {
|
||||
client = conn
|
||||
ClientReady = false
|
||||
Log.Info("client connected on ", zap.Any("address", client.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 Tcp {
|
||||
_, 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
|
||||
}
|
||||
if Pulsar {
|
||||
return PulsarSend(data)
|
||||
}
|
||||
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()
|
||||
}
|
||||
}
|
||||
-177
@@ -1,177 +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) {
|
||||
getWriteLock()
|
||||
db := getDb()
|
||||
defer db.Close()
|
||||
var insertSQL string
|
||||
if ClientReady && WriteToClient(teleString) == nil {
|
||||
insertSQL = InsertOld
|
||||
} else {
|
||||
insertSQL = InsertNew
|
||||
}
|
||||
stmt, _ := db.Prepare(insertSQL)
|
||||
|
||||
defer stmt.Close()
|
||||
result, err := stmt.Exec(teleString)
|
||||
id, _ := result.LastInsertId()
|
||||
checkDbErr(err, "error in insert telegram ")
|
||||
isWriting = false
|
||||
if insertSQL == InsertOld {
|
||||
Log.Info("telegram processed ", zap.Int64("id", id))
|
||||
} else {
|
||||
Log.Info("telegram saved ", zap.Int64("id", id))
|
||||
}
|
||||
}
|
||||
|
||||
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,45 +0,0 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"strings"
|
||||
|
||||
rotatelogs "github.com/lestrrat-go/file-rotatelogs"
|
||||
)
|
||||
|
||||
const (
|
||||
EndTag = "NNNN"
|
||||
Expression = "(?s)ZCZC.*NNNN"
|
||||
)
|
||||
|
||||
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"
|
||||
return check()
|
||||
}
|
||||
|
||||
func check() bool {
|
||||
teleString := buffer
|
||||
if strings.Index(teleString, EndTag) > 0 {
|
||||
telegram := exp.FindString(teleString)
|
||||
if len(telegram) > 0 {
|
||||
// Log.Info("telegram ", zap.String("text", telegram))
|
||||
telegram = removeEmpty(telegram) + "\n\n"
|
||||
telegram += "\n"
|
||||
_, _ = RawLog.Write([]byte(telegram + "\n"))
|
||||
buffer = ""
|
||||
InsertTelegram(telegram)
|
||||
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