add load unprocessed
This commit is contained in:
+1
-1
@@ -30,7 +30,7 @@ func Listen(address string) {
|
|||||||
if client == nil {
|
if client == nil {
|
||||||
client = conn
|
client = conn
|
||||||
log.Println("client connected ", client.RemoteAddr())
|
log.Println("client connected ", client.RemoteAddr())
|
||||||
|
LoadUnprocessed()
|
||||||
} else {
|
} else {
|
||||||
log.Println("client already connected, close")
|
log.Println("client already connected, close")
|
||||||
client.Close()
|
client.Close()
|
||||||
|
|||||||
+78
-21
@@ -21,8 +21,17 @@ const (
|
|||||||
`
|
`
|
||||||
InsertNew = "insert into telegram (tele_text) values (?)"
|
InsertNew = "insert into telegram (tele_text) values (?)"
|
||||||
InsertOld = "insert into telegram (tele_text, tele_processed) values (?, 1)"
|
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 (
|
var (
|
||||||
dbFile string
|
dbFile string
|
||||||
initTable bool
|
initTable bool
|
||||||
@@ -50,25 +59,14 @@ func InitDb(file string, init bool) error {
|
|||||||
db := getDb()
|
db := getDb()
|
||||||
defer db.Close()
|
defer db.Close()
|
||||||
_, dbErr = db.Exec(TableCreate)
|
_, dbErr = db.Exec(TableCreate)
|
||||||
if dbErr != nil {
|
checkDbErr(dbErr, "error in create table")
|
||||||
log.Fatal("error in create telegram table", dbErr)
|
|
||||||
}
|
|
||||||
log.Println("database ", db, " initialized")
|
log.Println("database ", db, " initialized")
|
||||||
initialized = true
|
initialized = true
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func InsertTelegram(teleString string) {
|
func InsertTelegram(teleString string) {
|
||||||
count := 0
|
getWriteLock()
|
||||||
for isWriting && count < 10 {
|
|
||||||
log.Println("waiting for write ")
|
|
||||||
time.Sleep(5 * time.Second)
|
|
||||||
count++
|
|
||||||
}
|
|
||||||
if count >= 5 {
|
|
||||||
log.Fatalln("database locked")
|
|
||||||
}
|
|
||||||
isWriting = true
|
|
||||||
db := getDb()
|
db := getDb()
|
||||||
defer db.Close()
|
defer db.Close()
|
||||||
var insertSQL string
|
var insertSQL string
|
||||||
@@ -77,17 +75,12 @@ func InsertTelegram(teleString string) {
|
|||||||
} else {
|
} else {
|
||||||
insertSQL = InsertNew
|
insertSQL = InsertNew
|
||||||
}
|
}
|
||||||
stmt, err := db.Prepare(insertSQL)
|
stmt, _ := db.Prepare(insertSQL)
|
||||||
if err != nil {
|
|
||||||
log.Fatal(err)
|
|
||||||
}
|
|
||||||
defer stmt.Close()
|
defer stmt.Close()
|
||||||
result, err := stmt.Exec(teleString)
|
result, err := stmt.Exec(teleString)
|
||||||
id, _ := result.LastInsertId()
|
id, _ := result.LastInsertId()
|
||||||
if err != nil {
|
checkDbErr(err, "error in insert telegram ")
|
||||||
isWriting = false
|
|
||||||
log.Fatalln("error in write telegram")
|
|
||||||
}
|
|
||||||
isWriting = false
|
isWriting = false
|
||||||
if insertSQL == InsertOld {
|
if insertSQL == InsertOld {
|
||||||
log.Println("telegram processed ", id)
|
log.Println("telegram processed ", id)
|
||||||
@@ -95,3 +88,67 @@ func InsertTelegram(teleString string) {
|
|||||||
log.Println("telegram saved ", id)
|
log.Println("telegram saved ", id)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func LoadUnprocessed() {
|
||||||
|
db := getDb()
|
||||||
|
defer db.Close()
|
||||||
|
for {
|
||||||
|
num := countTelegram()
|
||||||
|
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)
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user