package utils import ( "database/sql" "log" "os" "time" _ "github.com/mattn/go-sqlite3" ) const ( DbDriver = "sqlite3" TableCreate = ` 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 ) ` InsertNew = "insert into telegram (tele_text) values (?)" InsertOld = "insert into telegram (tele_text, tele_processed) values (?, 1)" 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=?" ) type telegram struct { id int64 text string } var ( dbFile string initTable bool initialized bool isWriting bool ) func getDb() *sql.DB { //log.Println("init database ", initTable) if initTable && !initialized { log.Println("remove database file", dbFile) os.Remove(dbFile) } db, err := sql.Open(DbDriver, dbFile) if err != nil { log.Fatal("error in open database ", err) } return db } func InitDb(file string, init bool) error { dbFile = file initTable = init var dbErr error db := getDb() defer db.Close() _, dbErr = db.Exec(TableCreate) checkDbErr(dbErr, "error in create table") log.Println("database ", db, " initialized") initialized = true return nil } func InsertTelegram(teleString string) { getWriteLock() db := getDb() defer db.Close() var insertSQL string if IsClientConnected() && WriteToClient(teleString) == nil { insertSQL = InsertOld } else { insertSQL = InsertNew } stmt, _ := db.Prepare(insertSQL) defer stmt.Close() result, err := stmt.Exec(teleString) id, _ := result.LastInsertId() checkDbErr(err, "error in insert telegram ") isWriting = false if insertSQL == InsertOld { log.Println("telegram processed ", id) } else { log.Println("telegram saved ", id) } } func LoadUnprocessed() { log.Println("loading telegram") db := getDb() defer db.Close() for { num := countTelegram() log.Println("unprocessed : ", num) 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 err := rows.Scan(&t.id, &t.text) checkDbErr(err, "error in scan telegram") if WriteToClient(t.text) == nil { processed(db, t) } } } } func processed(db *sql.DB, t telegram) { getWriteLock() tx, err := db.Begin() checkDbErr(err, "error in begin update") 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 log.Println("release lock") tx.Commit() } 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 log.Println("write lock :", isWriting) 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) } }