Files
tele-recv/app/app.go
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

268 lines
5.8 KiB
Go

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
}