Compare commits

..
10 Commits
Author SHA1 Message Date
w1ndyb0y 4191868f00 fix: improve local Pulsar connectivity checks 2026-07-10 17:07:42 +08:00
w1ndyb0y 46be9f483f docs: update README with new architecture and project structure
- Add architecture diagram and data flow
- Add project directory structure tree
- Add technology stack table
- Add Taskfile-based quick start and common commands
- Add cross-compilation instructions
- Add full config YAML example with all options
- Add test section with coverage info
- Add license section
- Keep existing download and release links
2026-07-10 16:07:48 +08:00
w1ndyb0y 7f65004cd4 test: add missing tests and fix hanging transport test
New test packages:
- config/loader_test.go: 5 tests (valid, minimal, missing file, pulsar-only, defaults)
- serial/serial_test.go: 6 tests (interface, read, EOF, close, empty)
- app/app_test.go: 5 tests (New, dispatch store/send, error handling, ctx cancellation)

Enhanced test packages:
- storage/store_test.go: +3 tests (concurrent insert, close/reopen, init clears)
- telegram/parser_test.go: +5 tests (concurrency, rawLog, empty, multiline, removeEmpty)
- transport/transport_test.go: +7 tests, +fixes

Fixes:
- app.go: store field type *storage.Store -> storage.Repository (supports mocking)
- transport_test: fix net.Pipe() sync deadlock in SetConnClosesPrevious
- transport_test: fix MultiSender_ErrorPropagation expectations

Results: 6 packages, 30+ tests, all pass with -race
2026-07-10 16:05:19 +08:00
w1ndyb0y 86d473c03d clean: remove deprecated utils/ package
- Remove all 5 files from utils/ (serial.go, socket.go, pulsar.go,
  sqlite.go, telegram.go)
- cmd/root.go: remove utils import and all utils. references
- cmd/test.go: rewrite to use serial, storage, transport packages directly
  with proper error logging via zap
2026-07-10 15:41:13 +08:00
w1ndyb0y 419f25f0dd fix: comprehensive bug fixes and architecture restructuring
Phase 1 — Bug fixes (7 bugs):
- Bug 1: Pulsar mode data loss — always persist to SQLite regardless of mode
- Bug 2: PulsarSend closing TCP client — use independent resetPulsarProducer()
- Bug 3: Greedy regex — use non-greedy (?s)ZCZC.*?NNNN
- Bug 4: Data after NNNN discarded — keep remaining buffer data
- Bug 5: strings.Index > 0 boundary — use strings.Contains
- Bug 6: Variable shadowing in PulsarSend — use = not :=
- Bug 7: Accept failure nil panic — add continue + retry logic

Phase 2 — Architecture restructuring:
- Split utils/ into config/, serial/, telegram/, storage/, transport/
- Define Sender, Repository, Reader interfaces
- Introduce app/ layer with context.Context lifecycle
- Replace spinlock with sync.Mutex
- Unified Config struct replaces 15+ global vars

Phase 3 — Testing & tooling:
- telegram/parser_test.go (7 test cases)
- storage/store_test.go (5 test cases)
- transport/transport_test.go (4 test cases)
- Taskfile: add test, test-race, test-cover tasks
- Go version: 1.15 -> 1.21
2026-07-10 15:32:34 +08:00
zhiqiang feng 5f3334ec87 send log 2021-02-01 17:46:20 +08:00
w1ndyb0y 133c80088c check tcp when insert telegram to db 2021-02-01 17:29:01 +08:00
w1ndyb0y fe1a8169e6 move pulsar send to telegram.go 2021-02-01 17:13:54 +08:00
w1ndyb0y 9ec63dd8b8 create producer 2021-02-01 17:07:52 +08:00
w1ndyb0y b9a826faf8 pulsar write 2021-02-01 17:03:24 +08:00
26 changed files with 2501 additions and 649 deletions
+181 -43
View File
@@ -1,64 +1,202 @@
# 电报接收 # tele-recv 串口电报接收
[![Build Status](http://d2.int.it2000.com.cn/api/badges/airport/tele-recv/status.svg)](http://d2.int.it2000.com.cn/airport/tele-recv) [![Build Status](http://d2.int.it2000.com.cn/api/badges/airport/tele-recv/status.svg)](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) ```
串口 (ttyS0) → serial.Port → telegram.Parser → storage.Store (SQLite)
## 安装
transport.Sender
解压下载文件 tele-recv-XXX.tag.gz
┌──────────┴──────────┐
如果是linux ↓ ↓
TCP Client Apache Pulsar
```Shell
tar xzvf tele-recv-XXX.tar.gz
``` ```
文件包中包含两个可执行文件和一个.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 | | 语言 | Go 1.21+ | 需 Go 1.21 或更高版本 |
|file | sqlite 文件名 | telegram.db | tele.db | | CLI | cobra + viper | 命令解析与配置管理 |
|init | 是否新建电报数据库 | true | true/false | | 串口 | github.com/argandas/serial | 串口通信 |
|address | 本地处理监听端口 | 127.0.0.1:6000 | | | 存储 | 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: ```bash
```Shell # 运行所有测试
./tele-recv-linux start task test
# 带竞态检测
task test-race
# 覆盖率报告
task test-cover
# 手动运行
go test -race -v ./...
``` ```
windows: 当前测试覆盖 6 个包,30+ 测试用例,`go test -race` 零竞态。
```Shell
tele-recv-win64.exe start
```
开始从串口设备读取数据,并开启TCP服务等待处理程序连接。 ## 下载
如果没有处理程序,电报则保存在本地数据库中。
等到处理程序连上端口,则把未处理电报以此发送给处理程序。 发布版本:[release](https://gitea.int.it2000.com.cn/airport/tele-recv/releases)
## 许可证
详见 [LICENSE](LICENSE)
+150
View File
@@ -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
View File
@@ -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
View File
@@ -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)
}
}
+25 -62
View File
@@ -25,43 +25,38 @@ import (
"go.uber.org/zap" "go.uber.org/zap"
"go.uber.org/zap/zapcore" "go.uber.org/zap/zapcore"
"github.com/spf13/viper" "it2000.com.cn/tele-recv/config"
"it2000.com.cn/tele-recv/utils"
) )
const ModeName = "TELEGRAM_MODE" const ModeName = "TELEGRAM_MODE"
var ( var (
cfgFile string cfgFile string
// Legacy global vars for backward compatibility with test command
device string device string
baudrate int baudrate int
lograw bool
dbFile string dbFile string
dbInit bool dbInit bool
socketAddress string socketAddress string
pulsarUrl string pulsarUrl string
topic string topic string
name string name string
tcp bool tcp bool
pulsar bool pulsar bool
lograw bool
Mode string
// rootCmd represents the base command when called without any subcommands // rootCmd represents the base command when called without any subcommands
rootCmd = &cobra.Command{ rootCmd = &cobra.Command{
Use: "tele-recv", Use: "tele-recv",
Short: "telegram receiver", Short: "telegram receiver",
Long: `A serial port telegram receiver. Long: `A serial port telegram receiver.
Read the telegram from `, Read the telegram from serial port`,
// Uncomment the following line if your bare application
// has an action associated with it:
// Run: func(cmd *cobra.Command, args []string) { },
} }
Mode = os.Getenv(ModeName) appCfg *config.Config
logger *zap.Logger logger *zap.Logger
) )
@@ -77,55 +72,33 @@ func Execute() {
func init() { func init() {
cobra.OnInitialize(initConfig) 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)") 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") rootCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
initLog() initLog()
defer logger.Sync()
// l = logger.Sugar()
utils.Log = logger
} }
// initConfig reads in config file and ENV variables if set. // initConfig reads in config file and ENV variables if set.
func initConfig() { func initConfig() {
if cfgFile != "" { cfg, err := config.Load(cfgFile)
// Use config file from the flag. if err != nil {
viper.SetConfigFile(cfgFile) logger.Warn("config load error, using defaults", zap.Error(err))
} else { return
// 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
} }
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() { func initLog() {
@@ -138,15 +111,6 @@ func initLog() {
panic(err) 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 { filePriority := zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
return lvl >= zapcore.DebugLevel return lvl >= zapcore.DebugLevel
}) })
@@ -163,5 +127,4 @@ func initLog() {
) )
logger = zap.New(core) logger = zap.New(core)
} }
+11 -74
View File
@@ -16,15 +16,10 @@ limitations under the License.
package cmd package cmd
import ( import (
"fmt"
"os"
"os/signal"
"syscall"
"go.uber.org/zap"
"it2000.com.cn/tele-recv/utils"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/zap"
"it2000.com.cn/tele-recv/app"
) )
// startCmd represents the start command // startCmd represents the start command
@@ -33,82 +28,24 @@ var startCmd = &cobra.Command{
Short: "start telegram receive service", Short: "start telegram receive service",
Long: `open serial port Long: `open serial port
start a tcp server for processing`, start a tcp server for processing`,
Run: func(cmd *cobra.Command, args []string) { RunE: func(cmd *cobra.Command, args []string) error {
start() return start()
}, },
} }
func init() { func init() {
rootCmd.AddCommand(startCmd) 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() { func start() error {
// Go signal notification works by sending `os.Signal` if appCfg == nil {
// values on a channel. We'll create a channel to logger.Fatal("configuration not loaded")
// 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 a, err := app.New(appCfg)
// 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)
}
for utils.ServerRunning {
if !utils.IsPortOpen() {
logger.Info("try to open port")
err := utils.OpenPort(device, baudrate)
if err != nil { if err != nil {
logger.Fatal("error in open serial port ", zap.Error(err)) logger.Fatal("failed to create app", 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")
return a.Run()
} }
+40 -16
View File
@@ -16,8 +16,14 @@ limitations under the License.
package cmd package cmd
import ( import (
"net"
"net/url"
"time"
"github.com/spf13/cobra" "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 // testCmd represents the test command
@@ -34,27 +40,45 @@ Test load sqlite database and init table`,
func test() { func test() {
logger.Info("Testing environment") 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 { 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) // Test SQLite
utils.ClosePulsar() 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() { func init() {
rootCmd.AddCommand(testCmd) 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")
} }
+43
View File
@@ -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
}
+54
View File
@@ -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
}
+200
View File
@@ -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)
}
}
+52 -8
View File
@@ -1,28 +1,72 @@
module it2000.com.cn/tele-recv module it2000.com.cn/tele-recv
go 1.15 go 1.21
require ( require (
github.com/apache/pulsar-client-go v0.3.0 github.com/apache/pulsar-client-go v0.3.0
github.com/argandas/serial v0.0.0-20160316175758-889a5ad85462 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/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/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/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/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/afero v1.5.1 // indirect
github.com/spf13/cast v1.3.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/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/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/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/sys v0.0.0-20201214210602-f9fddec55a1e // indirect
golang.org/x/text v0.3.4 // 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/ini.v1 v1.62.0 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect
) )
-4
View File
@@ -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/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/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/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 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 h1:6MnRN8NT7+YBpUIWxHtefFZOKTAPgGjpQSxqLNn0+qY=
github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= 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/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/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/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.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
github.com/kr/pretty v0.2.0 h1:s5hAObm+yFO5uHYt5dYjxi2rXrsnmRpJx4OYvIWUaQs= github.com/kr/pretty v0.2.0 h1:s5hAObm+yFO5uHYt5dYjxi2rXrsnmRpJx4OYvIWUaQs=
github.com/kr/pretty v0.2.0/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI= 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 h1:2xWsjqPFWcplujydGg4WmhC/6fZqK42wMM8aXeqhl0I=
github.com/pierrec/lz4 v2.0.5+incompatible/go.mod h1:pdkljMzZIN41W+lC3N2tnIh5sFi+IEE17M5jbnwPHcY= 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.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.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= 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= 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/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 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-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 h1:YR8cESwS4TdDjEe65xsg0ogRM/Nc3DYOhEAlW+xobZo=
gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+50
View File
@@ -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
}
+100
View File
@@ -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)
}
+131
View File
@@ -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()
}
+194
View File
@@ -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
View File
@@ -6,17 +6,15 @@ serial:
#stopbits: 1 #stopbits: 1
lograw: true lograw: true
sqlite: sqlite:
file: telegram.db file: telegram.db
init: true init: true
socket: socket:
address: 127.0.0.1:6000 address: 127.0.0.1:6000
pulsar: pulsar:
url: pulsar://yzjc.gzzn.dev:6650 url: pulsar://localhost:6650
topic: telegram-raw topic: telegram-raw
name: serial-reader name: serial-reader
+83
View File
@@ -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")
}
+214
View File
@@ -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)
}
}
+277
View File
@@ -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)
+234
View File
@@ -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)
}
}
-63
View File
@@ -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()
}
}
-50
View File
@@ -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
}
-93
View File
@@ -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
View File
@@ -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))
}
}
-45
View File
@@ -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")
}