Merge branch 'develop' of gitee.com:w1ndyb0y/tele-recv into develop

 Conflicts:
	cmd/start.go
This commit is contained in:
w1ndyb0y
2020-12-23 11:13:36 +08:00
5 changed files with 89 additions and 39 deletions
+11 -11
View File
@@ -17,26 +17,25 @@ package cmd
import ( import (
"fmt" "fmt"
"it2000.com.cn/tele-recv/utils"
"log" "log"
"os" "os"
"os/signal" "os/signal"
"syscall" "syscall"
"it2000.com.cn/tele-recv/utils"
"github.com/spf13/cobra" "github.com/spf13/cobra"
) )
// startCmd represents the start command // startCmd represents the start command
var startCmd = &cobra.Command{ var startCmd = &cobra.Command{
Use: "start", Use: "start",
Short: "A brief description of your command", Short: "start telegram receive service",
Long: `A longer description that spans multiple lines and likely contains examples Long: `open serial port
and usage of using your command. For example: start a tcp server for processing`,
Run: func(cmd *cobra.Command, args []string) {
Cobra is a CLI library for Go that empowers applications. start()
This application is a tool to generate the needed files },
to quickly create a Cobra application.`,
Run: start,
} }
func init() { func init() {
@@ -53,7 +52,7 @@ func init() {
// startCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle") // startCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
} }
func start(cmd *cobra.Command, args []string) { func start() {
// Go signal notification works by sending `os.Signal` // Go signal notification works by sending `os.Signal`
// values on a channel. We'll create a channel to // values on a channel. We'll create a channel to
// receive these notifications (we'll also make one to // receive these notifications (we'll also make one to
@@ -74,6 +73,7 @@ func start(cmd *cobra.Command, args []string) {
fmt.Println() fmt.Println()
fmt.Println(sig) fmt.Println(sig)
utils.ServerRunning = false utils.ServerRunning = false
utils.StopSocketServer()
done <- true done <- true
}() }()
@@ -82,7 +82,7 @@ func start(cmd *cobra.Command, args []string) {
// above sending a value on `done`) and then exit. // above sending a value on `done`) and then exit.
fmt.Println("awaiting signal") fmt.Println("awaiting signal")
_ = utils.InitDb(dbFile, dbInit) _ = utils.InitDb(dbFile, dbInit)
utils.Listen(socketAddress) go utils.Listen(socketAddress)
for utils.ServerRunning { for utils.ServerRunning {
if !utils.IsPortOpen() { if !utils.IsPortOpen() {
log.Println("try to open port") log.Println("try to open port")
+4 -2
View File
@@ -28,10 +28,12 @@ var testCmd = &cobra.Command{
Long: `loading config file Long: `loading config file
Test serial port Test serial port
Test load sqlite database and init table`, Test load sqlite database and init table`,
Run: test, Run: func(cmd *cobra.Command, args []string) {
test()
},
} }
func test(cmd *cobra.Command, args []string) { func test() {
log.Println("Testing environment") log.Println("Testing environment")
_ = utils.OpenPort(device, baudrate, bytesize) _ = utils.OpenPort(device, baudrate, bytesize)
+11 -8
View File
@@ -1,26 +1,27 @@
package utils package utils
import ( import (
"github.com/tarm/serial"
"io" "io"
"log" "log"
"time" "time"
"github.com/tarm/serial"
) )
var( var (
device string device string
baudrate int baudrate int
bytesize byte bytesize byte
port *serial.Port port *serial.Port
portOpened bool portOpened = false
) )
func OpenPort(device string, baudrate int, bytesize byte) error { func OpenPort(device string, baudrate int, bytesize byte) error {
log.Println("opening port device:", device, ", baudrate=", baudrate, ", bytesize=",bytesize) log.Println("opening port device:", device, ", baudrate=", baudrate, ", bytesize=", bytesize)
config := &serial.Config{Name: device, Baud: baudrate, Size: bytesize, ReadTimeout: time.Second * 1} config := &serial.Config{Name: device, Baud: baudrate, Size: bytesize, ReadTimeout: time.Second * 1}
var err error var err error
port, err = serial.OpenPort(config) port, err = serial.OpenPort(config)
if err!=nil { if err != nil {
log.Fatalln("error in open serial port ", err) log.Fatalln("error in open serial port ", err)
return err return err
} }
@@ -29,9 +30,12 @@ func OpenPort(device string, baudrate int, bytesize byte) error {
return nil return nil
} }
func ReadPort() ([]byte, int){ func ReadPort() ([]byte, int) {
buffer := make([]byte, 256) buffer := make([]byte, 256)
count,err := port.Read(buffer) if !ServerRunning {
return buffer, 0
}
count, err := port.Read(buffer)
if err != nil && err != io.EOF { if err != nil && err != io.EOF {
log.Println("error in read port ", err) log.Println("error in read port ", err)
//TODO: 3 time fail try to //TODO: 3 time fail try to
@@ -42,4 +46,3 @@ func ReadPort() ([]byte, int){
func IsPortOpen() bool { func IsPortOpen() bool {
return portOpened return portOpened
} }
+40 -7
View File
@@ -1,17 +1,21 @@
package utils package utils
import ( import (
"io"
"log" "log"
"net" "net"
"time"
) )
const ServerType = "tcp" const ServerType = "tcp"
var ServerRunning = false var ServerRunning = false
var client net.Conn var client net.Conn
var server net.Listener
func Listen(address string) { func Listen(address string) {
server, err := net.Listen(ServerType, address) var err error
server, err = net.Listen(ServerType, address)
if err != nil { if err != nil {
log.Fatal("error in create server ", err) log.Fatal("error in create server ", err)
} }
@@ -25,22 +29,51 @@ func Listen(address string) {
} }
if client == nil { if client == nil {
client = conn client = conn
log.Println("client connected ", client.RemoteAddr())
} else { } else {
log.Println("client already connected, close") log.Println("client already connected, close")
client.Close()
client = nil
conn.Close() conn.Close()
} }
if IsClientConnected() {
connCheck()
}
time.Sleep(2 * time.Second)
} }
} }
func Write(data string) { func WriteToClient(data string) error {
defer func() { _, err := client.Write([]byte(data + "\r\n"))
client.Close()
client = nil
}()
_, err := client.Write([]byte(data))
if err != nil { if err != nil {
log.Println("error in write data to client socket ", err) log.Println("error in write data to client socket ", err)
client.Close() client.Close()
client = nil client = nil
} }
return err
}
func connCheck() bool {
_, err := client.Read(make([]byte, 0))
if err != nil && err != io.EOF {
// this connection is invalid
log.Println("conn closed....", err)
client = nil
return false
}
return true
}
func IsClientConnected() bool {
return client != nil
}
func StopSocketServer() {
if IsClientConnected() {
client.Close()
}
if server != nil {
server.Close()
}
} }
+18 -6
View File
@@ -2,10 +2,11 @@ package utils
import ( import (
"database/sql" "database/sql"
_ "github.com/mattn/go-sqlite3"
"log" "log"
"os" "os"
"time" "time"
_ "github.com/mattn/go-sqlite3"
) )
const ( const (
@@ -18,7 +19,8 @@ const (
[tele_text] TEXT NOT NULL [tele_text] TEXT NOT NULL
) )
` `
TableInsert = "insert into telegram (tele_text) values (?)" InsertNew = "insert into telegram (tele_text) values (?)"
InsertOld = "insert into telegram (tele_text, tele_processed) values (?, 1)"
) )
var ( var (
@@ -58,9 +60,9 @@ func InitDb(file string, init bool) error {
func InsertTelegram(teleString string) { func InsertTelegram(teleString string) {
count := 0 count := 0
for isWriting && count < 5 { for isWriting && count < 10 {
log.Println("waiting for write ") log.Println("waiting for write ")
time.Sleep(3 * time.Second) time.Sleep(5 * time.Second)
count++ count++
} }
if count >= 5 { if count >= 5 {
@@ -69,7 +71,13 @@ func InsertTelegram(teleString string) {
isWriting = true isWriting = true
db := getDb() db := getDb()
defer db.Close() defer db.Close()
stmt, err := db.Prepare(TableInsert) var insertSQL string
if IsClientConnected() && WriteToClient(teleString) == nil {
insertSQL = InsertOld
} else {
insertSQL = InsertNew
}
stmt, err := db.Prepare(insertSQL)
if err != nil { if err != nil {
log.Fatal(err) log.Fatal(err)
} }
@@ -80,6 +88,10 @@ func InsertTelegram(teleString string) {
isWriting = false isWriting = false
log.Fatalln("error in write telegram") log.Fatalln("error in write telegram")
} }
log.Println("telegram insert ", id)
isWriting = false isWriting = false
if insertSQL == InsertOld {
log.Println("telegram processed ", id)
} else {
log.Println("telegram saved ", id)
}
} }