Files
tele-recv/utils/sqlite.go
T

157 lines
3.4 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"
"log"
"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-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 {
log.Println("remove database file", dbFile)
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 {
log.Fatal("error in open database ", err)
}
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-21 17:58:55 +08:00
log.Println("database ", db, " initialized")
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
if IsClientConnected() && WriteToClient(teleString) == nil {
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 {
log.Println("telegram processed ", id)
} else {
log.Println("telegram saved ", id)
}
2020-12-21 17:58:55 +08:00
}
2020-12-23 13:07:45 +08:00
func LoadUnprocessed() {
2020-12-23 13:41:46 +08:00
log.Println("loading telegram")
2020-12-23 13:07:45 +08:00
db := getDb()
defer db.Close()
for {
num := countTelegram()
2020-12-23 13:41:46 +08:00
log.Println("unprocessed : ", num)
2020-12-23 13:07:45 +08:00
if num < 1 {
break
}
log.Println("try to load unprocessed telegram ", num)
rows, err := db.Query(load)
checkDbErr(err, "error in query telegram")
defer rows.Close()
for rows.Next() {
var t telegram
2020-12-23 13:41:46 +08:00
err := rows.Scan(&t.id, &t.text)
2020-12-23 13:07:45 +08:00
checkDbErr(err, "error in scan telegram")
if WriteToClient(t.text) == nil {
processed(t)
}
}
}
}
func processed(t telegram) {
db := getDb()
getWriteLock()
stmt, _ := db.Prepare(update)
result, err := stmt.Exec(t.id)
checkDbErr(err, "error in update telegram status")
id, _ := result.RowsAffected()
log.Println(t.id, " telegram update : ", id)
isWriting = false
}
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
for isWriting && count < 10 {
log.Println("waiting for write ")
time.Sleep(5 * time.Second)
count++
}
if count >= 10 {
log.Fatalln("database locked")
}
isWriting = true
}
func checkDbErr(err error, msg string) {
if err != nil {
log.Fatalln(msg, err)
}
}