Files
tele-recv/utils/sqlite.go
T

184 lines
4.1 KiB
Go
Raw Normal View History

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
2021-02-01 17:03:24 +08:00
if Pulsar {
err := PulsarSend(teleString)
if err != nil {
Log.Error("error in send to pulsar", zap.Error(err))
}
}
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
if insertSQL == InsertOld {
2020-12-28 15:58:00 +08:00
Log.Info("telegram processed ", zap.Int64("id", id))
} else {
2020-12-28 15:58:00 +08:00
Log.Info("telegram saved ", zap.Int64("id", id))
}
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-28 16:09:08 +08:00
time.Sleep(500 * time.Millisecond)
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")
2020-12-28 16:09:08 +08:00
rows, _ := result.RowsAffected()
2020-12-28 17:16:38 +08:00
Log.Info("telegram update ", zap.Int64("id", t.id), zap.Bool("result", rows > 0))
2020-12-23 13:07:45 +08:00
isWriting = false
2020-12-28 17:16:38 +08:00
Log.Debug("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
}
}