2020-12-21 15:38:33 +08:00
|
|
|
package utils
|
2020-12-21 11:19:53 +08:00
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"database/sql"
|
|
|
|
|
"os"
|
2020-12-21 17:58:55 +08:00
|
|
|
"time"
|
2020-12-22 11:45:25 +08:00
|
|
|
|
|
|
|
|
_ "github.com/mattn/go-sqlite3"
|
2020-12-28 15:58:00 +08:00
|
|
|
"go.uber.org/zap"
|
2020-12-21 11:19:53 +08:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
const (
|
|
|
|
|
DbDriver = "sqlite3"
|
2020-12-21 17:58:55 +08:00
|
|
|
TableCreate = `
|
2020-12-21 11:19:53 +08:00
|
|
|
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
|
|
|
|
|
)
|
|
|
|
|
`
|
2020-12-22 11:45:25 +08:00
|
|
|
InsertNew = "insert into telegram (tele_text) values (?)"
|
|
|
|
|
InsertOld = "insert into telegram (tele_text, tele_processed) values (?, 1)"
|
2020-12-23 13:07:45 +08:00
|
|
|
|
|
|
|
|
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=?"
|
2020-12-21 11:19:53 +08:00
|
|
|
)
|
|
|
|
|
|
2020-12-23 13:07:45 +08:00
|
|
|
type telegram struct {
|
|
|
|
|
id int64
|
|
|
|
|
text string
|
|
|
|
|
}
|
|
|
|
|
|
2020-12-21 17:58:55 +08:00
|
|
|
var (
|
|
|
|
|
dbFile string
|
|
|
|
|
initTable bool
|
|
|
|
|
initialized bool
|
|
|
|
|
isWriting bool
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func getDb() *sql.DB {
|
2020-12-23 13:45:00 +08:00
|
|
|
//log.Println("init database ", initTable)
|
2020-12-21 17:58:55 +08:00
|
|
|
if initTable && !initialized {
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("remove database ", zap.String("file", dbFile))
|
2020-12-21 17:58:55 +08:00
|
|
|
os.Remove(dbFile)
|
2020-12-21 11:19:53 +08:00
|
|
|
}
|
2020-12-21 17:58:55 +08:00
|
|
|
db, err := sql.Open(DbDriver, dbFile)
|
2020-12-21 11:19:53 +08:00
|
|
|
if err != nil {
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Fatal("error in open database ", zap.Error(err))
|
2020-12-21 11:19:53 +08:00
|
|
|
}
|
2020-12-21 17:58:55 +08:00
|
|
|
return db
|
2020-12-21 11:36:08 +08:00
|
|
|
}
|
|
|
|
|
|
2020-12-21 15:38:33 +08:00
|
|
|
func InitDb(file string, init bool) error {
|
2020-12-21 17:58:55 +08:00
|
|
|
dbFile = file
|
|
|
|
|
initTable = init
|
2020-12-21 11:36:08 +08:00
|
|
|
var dbErr error
|
2020-12-21 17:58:55 +08:00
|
|
|
db := getDb()
|
|
|
|
|
defer db.Close()
|
|
|
|
|
_, dbErr = db.Exec(TableCreate)
|
2020-12-23 13:07:45 +08:00
|
|
|
checkDbErr(dbErr, "error in create table")
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("database initialized")
|
2020-12-21 17:58:55 +08:00
|
|
|
initialized = true
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func InsertTelegram(teleString string) {
|
2020-12-23 13:07:45 +08:00
|
|
|
getWriteLock()
|
2020-12-21 17:58:55 +08:00
|
|
|
db := getDb()
|
|
|
|
|
defer db.Close()
|
2020-12-22 11:45:25 +08:00
|
|
|
var insertSQL string
|
2020-12-23 14:30:09 +08:00
|
|
|
if ClientReady && WriteToClient(teleString) == nil {
|
2020-12-22 11:45:25 +08:00
|
|
|
insertSQL = InsertOld
|
|
|
|
|
} else {
|
|
|
|
|
insertSQL = InsertNew
|
|
|
|
|
}
|
2020-12-23 13:07:45 +08:00
|
|
|
stmt, _ := db.Prepare(insertSQL)
|
|
|
|
|
|
2020-12-21 17:58:55 +08:00
|
|
|
defer stmt.Close()
|
|
|
|
|
result, err := stmt.Exec(teleString)
|
|
|
|
|
id, _ := result.LastInsertId()
|
2020-12-23 13:07:45 +08:00
|
|
|
checkDbErr(err, "error in insert telegram ")
|
2020-12-21 17:58:55 +08:00
|
|
|
isWriting = false
|
2020-12-22 12:09:33 +08:00
|
|
|
if insertSQL == InsertOld {
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("telegram processed ", zap.Int64("id", id))
|
2020-12-22 12:09:33 +08:00
|
|
|
} else {
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("telegram saved ", zap.Int64("id", id))
|
2020-12-22 12:09:33 +08:00
|
|
|
}
|
2020-12-21 17:58:55 +08:00
|
|
|
}
|
2020-12-23 13:07:45 +08:00
|
|
|
|
|
|
|
|
func LoadUnprocessed() {
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Info("loading telegram")
|
2020-12-23 13:07:45 +08:00
|
|
|
db := getDb()
|
|
|
|
|
defer db.Close()
|
|
|
|
|
for {
|
2020-12-23 14:23:55 +08:00
|
|
|
telegrams := getTelegram(db)
|
|
|
|
|
if len(telegrams) < 1 {
|
2020-12-23 13:07:45 +08:00
|
|
|
break
|
|
|
|
|
}
|
2020-12-23 14:23:55 +08:00
|
|
|
for _, t := range telegrams {
|
2020-12-23 13:07:45 +08:00
|
|
|
if WriteToClient(t.text) == nil {
|
2020-12-23 14:02:16 +08:00
|
|
|
processed(db, t)
|
2020-12-23 13:07:45 +08:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2020-12-23 14:30:09 +08:00
|
|
|
ClientReady = true
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Info("client is ready for new telegram")
|
2020-12-23 13:07:45 +08:00
|
|
|
}
|
|
|
|
|
|
2020-12-23 14:23:55 +08:00
|
|
|
func getTelegram(db *sql.DB) []telegram {
|
|
|
|
|
var telegrams []telegram
|
|
|
|
|
num := countTelegram()
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("telegram ", zap.Int64("unprocessed", num))
|
2020-12-23 14:23:55 +08:00
|
|
|
if num < 1 {
|
|
|
|
|
return telegrams
|
|
|
|
|
}
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Info("try to load unprocessed telegram ")
|
2020-12-23 14:23:55 +08:00
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
|
2020-12-23 14:02:16 +08:00
|
|
|
func processed(db *sql.DB, t telegram) {
|
2020-12-23 13:07:45 +08:00
|
|
|
getWriteLock()
|
2020-12-23 14:04:06 +08:00
|
|
|
tx, err := db.Begin()
|
|
|
|
|
checkDbErr(err, "error in begin update")
|
2020-12-23 14:12:00 +08:00
|
|
|
stmt, err := db.Prepare(update)
|
|
|
|
|
defer stmt.Close()
|
|
|
|
|
checkDbErr(err, "error in create update statement")
|
2020-12-23 13:07:45 +08:00
|
|
|
result, err := stmt.Exec(t.id)
|
|
|
|
|
checkDbErr(err, "error in update telegram status")
|
|
|
|
|
id, _ := result.RowsAffected()
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Info("telegram update : %t", zap.Int64("id", id), zap.Bool("result", id > 0))
|
2020-12-23 13:07:45 +08:00
|
|
|
isWriting = false
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Info("release lock")
|
2020-12-23 14:04:06 +08:00
|
|
|
tx.Commit()
|
2020-12-23 13:07:45 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Debug("with ", zap.Bool("lock", isWriting))
|
2020-12-23 13:07:45 +08:00
|
|
|
for isWriting && count < 10 {
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Info("waiting for write ")
|
2020-12-23 13:07:45 +08:00
|
|
|
time.Sleep(5 * time.Second)
|
|
|
|
|
count++
|
|
|
|
|
}
|
|
|
|
|
if count >= 10 {
|
2020-12-28 11:06:36 +08:00
|
|
|
Log.Fatal("database locked")
|
2020-12-23 13:07:45 +08:00
|
|
|
}
|
|
|
|
|
isWriting = true
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func checkDbErr(err error, msg string) {
|
|
|
|
|
if err != nil {
|
2020-12-28 15:58:00 +08:00
|
|
|
Log.Fatal(msg, zap.Error(err))
|
2020-12-23 13:07:45 +08:00
|
|
|
}
|
|
|
|
|
}
|