package utils import ( "database/sql" "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.Info("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.Infof("database %s initialized", dbFile) initialized = true return nil } func InsertTelegram(teleString string) { getWriteLock() db := getDb() defer db.Close() var insertSQL string if ClientReady && 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.Infof("telegram [%d] processed ", id) } else { Log.Infof("telegram [%d] saved ", id) } } func LoadUnprocessed() { Log.Info("loading telegram") db := getDb() defer db.Close() for { telegrams := getTelegram(db) if len(telegrams) < 1 { break } for _, t := range telegrams { if WriteToClient(t.text) == nil { processed(db, t) } } } ClientReady = true Log.Info("client is ready for new telegram") } func getTelegram(db *sql.DB) []telegram { var telegrams []telegram num := countTelegram() Log.Infof("%n telegram unprocessed ", num) if num < 1 { return telegrams } Log.Info("try to load unprocessed telegram ") 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 } func processed(db *sql.DB, t telegram) { getWriteLock() tx, err := db.Begin() checkDbErr(err, "error in begin update") stmt, err := db.Prepare(update) defer stmt.Close() checkDbErr(err, "error in create update statement") result, err := stmt.Exec(t.id) checkDbErr(err, "error in update telegram status") id, _ := result.RowsAffected() Log.Infof("[%d] telegram update : %t", id, id > 0) isWriting = false Log.Info("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.Debugf("write lock : %t", isWriting) for isWriting && count < 10 { Log.Info("waiting for write ") time.Sleep(5 * time.Second) count++ } if count >= 10 { Log.Fatal("database locked") } isWriting = true } func checkDbErr(err error, msg string) { if err != nil { Log.Fatal(msg, err) } }