package db import ( "database/sql" "encoding/json" "fmt" "os" "path/filepath" "sync" "time" _ "modernc.org/sqlite" ) var ( dbPath string mu sync.Mutex sqlDB *sql.DB ) type Room struct { LiveID string `json:"live_id"` Keywords string `json:"keywords"` CreatedAt float64 `json:"created_at"` } type User struct { ID int64 `json:"id"` LiveID string `json:"live_id"` UserID string `json:"user_id"` UserName string `json:"user_name"` FirstSeen float64 `json:"first_seen"` LastSeen float64 `json:"last_seen"` MsgCount int `json:"msg_count"` GiftCount int `json:"gift_count"` EnterCount int `json:"enter_count"` LikeCount int `json:"like_count"` FollowCount int `json:"follow_count"` LastMsg string `json:"last_msg"` LastAction string `json:"last_action"` HitKw string `json:"hit_kw"` SecUID string `json:"sec_uid"` Avatar string `json:"avatar"` Mark string `json:"mark"` Lead int `json:"lead"` LeadTs float64 `json:"lead_ts"` } type ChatLog struct { ID int64 `json:"id"` LiveID string `json:"live_id"` Msg string `json:"msg"` Seen int `json:"seen"` } type ErrorLog struct { ID int64 `json:"id"` TS float64 `json:"ts"` Source string `json:"source"` RoomID string `json:"room_id"` Message string `json:"message"` Stack string `json:"stack"` } func init() { if exePath, err := os.Executable(); err == nil { candidate := filepath.Join(filepath.Dir(exePath), "data.db") if dir := filepath.Dir(candidate); writable(dir) { dbPath = candidate return } } appData := os.Getenv("APPDATA") if appData == "" { appData = filepath.Join(os.Getenv("USERPROFILE"), "AppData", "Roaming") } dbPath = filepath.Join(appData, "青梧映星", "douyinlive", "data.db") } func writable(d string) bool { if err := os.MkdirAll(d, 0755); err != nil { return false } f, err := os.CreateTemp(d, ".wtest") if err != nil { return false } f.Close() os.Remove(f.Name()) return true } func conn() (*sql.DB, error) { if sqlDB != nil { return sqlDB, nil } dir := filepath.Dir(dbPath) if err := os.MkdirAll(dir, 0755); err != nil { return nil, err } var err error sqlDB, err = sql.Open("sqlite", dbPath) if err != nil { return nil, err } sqlDB.SetMaxOpenConns(4) sqlDB.SetMaxIdleConns(2) return sqlDB, nil } func InitDB() error { mu.Lock() defer mu.Unlock() c, err := conn() if err != nil { return err } return ensureMigrated(c) } func ensureMigrated(c *sql.DB) error { _, err := c.Exec(` CREATE TABLE IF NOT EXISTS rooms(live_id TEXT PRIMARY KEY, keywords TEXT DEFAULT NULL, created_at REAL); CREATE TABLE IF NOT EXISTS users(id INTEGER PRIMARY KEY AUTOINCREMENT, live_id TEXT, user_id TEXT, user_name TEXT, first_seen REAL, last_seen REAL, msg_count INTEGER DEFAULT 0, gift_count INTEGER DEFAULT 0, enter_count INTEGER DEFAULT 0, like_count INTEGER DEFAULT 0, follow_count INTEGER DEFAULT 0, last_msg TEXT, last_action TEXT, hit_kw TEXT, sec_uid TEXT DEFAULT '', avatar TEXT DEFAULT '', mark TEXT DEFAULT '', lead INTEGER DEFAULT 0, lead_ts REAL DEFAULT 0, UNIQUE(live_id, user_id)); CREATE TABLE IF NOT EXISTS error_logs(id INTEGER PRIMARY KEY AUTOINCREMENT, ts REAL, source TEXT, room_id TEXT, message TEXT, stack TEXT); CREATE INDEX IF NOT EXISTS idx_el ON error_logs(ts); CREATE TABLE IF NOT EXISTS settings(k TEXT PRIMARY KEY, v TEXT); CREATE TABLE IF NOT EXISTS chat_logs(id INTEGER PRIMARY KEY AUTOINCREMENT, live_id TEXT, ts REAL, msg TEXT, seen INTEGER DEFAULT 0); CREATE INDEX IF NOT EXISTS idx_cl ON chat_logs(live_id, id); CREATE INDEX IF NOT EXISTS idx_ul ON users(live_id); CREATE INDEX IF NOT EXISTS idx_un ON users(user_name); `) if err != nil { return err } return migrateUsers(c) } func migrateUsers(c *sql.DB) error { var cnt int c.QueryRow("SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='users'").Scan(&cnt) if cnt == 0 { return nil } for col, ddl := range map[string]string{ "sec_uid": "TEXT DEFAULT ''", "avatar": "TEXT DEFAULT ''", "mark": "TEXT DEFAULT ''", "lead": "INTEGER DEFAULT 0", "lead_ts": "REAL DEFAULT 0", } { var n int c.QueryRow("SELECT COUNT(*) FROM pragma_table_info('users') WHERE name=?", col).Scan(&n) if n == 0 { c.Exec(fmt.Sprintf("ALTER TABLE users ADD COLUMN %s %s", col, ddl)) } } var n int c.QueryRow("SELECT COUNT(*) FROM pragma_table_info('chat_logs') WHERE name='seen'").Scan(&n) if n == 0 { c.Exec("ALTER TABLE chat_logs ADD COLUMN seen INTEGER DEFAULT 0") } return nil } // SaveRoom records a room (idempotent on live_id). func SaveRoom(liveID, keywords string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("INSERT OR IGNORE INTO rooms(live_id, keywords, created_at) VALUES(?,?,?)", liveID, keywords, float64(time.Now().Unix())) } // GetRooms returns all rooms from DB. func GetRooms() []Room { mu.Lock() defer mu.Unlock() c, _ := conn() rows, err := c.Query("SELECT live_id, keywords, created_at FROM rooms ORDER BY created_at ASC") if err != nil { return nil } defer rows.Close() var out []Room for rows.Next() { var r Room rows.Scan(&r.LiveID, &r.Keywords, &r.CreatedAt) out = append(out, r) } return out } // DelRoom removes a room and its associated users/chat_logs. func DelRoom(liveID string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("DELETE FROM users WHERE live_id=?", liveID) c.Exec("DELETE FROM rooms WHERE live_id=?", liveID) c.Exec("DELETE FROM chat_logs WHERE live_id=?", liveID) } // GetUsers returns users for a room. func GetUsers(liveID string) []User { mu.Lock() defer mu.Unlock() c, _ := conn() rows, err := c.Query( `SELECT id,live_id,user_id,user_name,first_seen,last_seen, msg_count,gift_count,enter_count,like_count,follow_count, last_msg,last_action,hit_kw,sec_uid,avatar,mark,lead,lead_ts FROM users WHERE live_id=? ORDER BY (msg_count*10+gift_count*20+enter_count+like_count+follow_count) DESC`, liveID) if err != nil { return nil } defer rows.Close() var out []User for rows.Next() { var u User rows.Scan(&u.ID, &u.LiveID, &u.UserID, &u.UserName, &u.FirstSeen, &u.LastSeen, &u.MsgCount, &u.GiftCount, &u.EnterCount, &u.LikeCount, &u.FollowCount, &u.LastMsg, &u.LastAction, &u.HitKw, &u.SecUID, &u.Avatar, &u.Mark, &u.Lead, &u.LeadTs) out = append(out, u) } return out } // UpsertUser creates or updates a user row. func UpsertUser(liveID, uid, uname, action string, extra map[string]string) { mu.Lock() defer mu.Unlock() c, _ := conn() now := float64(time.Now().Unix()) var id int64 var mc, gc, ec, lc, fc int err := c.QueryRow("SELECT id,msg_count,gift_count,enter_count,like_count,follow_count FROM users WHERE live_id=? AND user_id=?", liveID, uid).Scan(&id, &mc, &gc, &ec, &lc, &fc) if err == nil { switch action { case "chat": mc++ case "gift": gc++ case "enter": ec++ case "like": lc++ case "follow": fc++ } q := "'`" sets := []string{fmt.Sprintf("last_seen=%g", now), fmt.Sprintf("last_action=%s%s%s", q, action, q)} if action == "chat" { cm := escapeSQL(extra["content"]) kw := escapeSQL(extra["keyword"]) sets = append(sets, fmt.Sprintf("msg_count=%d", mc)) sets = append(sets, fmt.Sprintf("last_msg=%s%s%s", q, cm, q)) sets = append(sets, fmt.Sprintf("hit_kw=%s%s%s", q, kw, q)) } else if action == "gift" { sets = append(sets, fmt.Sprintf("gift_count=%d", gc)) } else if action == "enter" { sets = append(sets, fmt.Sprintf("enter_count=%d", ec)) } else if action == "like" { sets = append(sets, fmt.Sprintf("like_count=%d", lc)) } else if action == "follow" { sets = append(sets, fmt.Sprintf("follow_count=%d", fc)) } for _, col := range []string{"sec_uid", "avatar"} { if v := extra[col]; v != "" { sets = append(sets, fmt.Sprintf("%s=%s%s%s", col, q, escapeSQL(v), q)) } } c.Exec(fmt.Sprintf("UPDATE users SET %s WHERE id=%d", joinStrings(sets, ","), id)) } else { cols := []string{"live_id", "user_id", "user_name", "first_seen", "last_seen", "last_action"} vs := []interface{}{liveID, uid, uname, now, now, action} for _, col := range []string{"sec_uid", "avatar"} { if v := extra[col]; v != "" { cols = append(cols, col) vs = append(vs, v) } } ps := make([]string, len(cols)) for i := range ps { ps[i] = "?" } c.Exec(fmt.Sprintf("INSERT INTO users(%s) VALUES(%s)", joinStrings(cols, ","), joinStrings(ps, ",")), vs...) } } func escapeSQL(s string) string { return `''` + s + `''` } func joinStrings(ss []string, sep string) string { if len(ss) == 0 { return "" } out := ss[0] for i := 1; i < len(ss); i++ { out += sep + ss[i] } return out } // SetMark assigns a color/mark to a user. func SetMark(liveID, userID, mark string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("UPDATE users SET mark=? WHERE live_id=? AND user_id=?", mark, liveID, userID) if mark != "" { var id int64 c.QueryRow("SELECT id FROM users WHERE live_id=? AND user_id=?", liveID, userID).Scan(&id) if id == 0 { now := float64(time.Now().Unix()) c.Exec("INSERT OR IGNORE INTO users(live_id,user_id,user_name,first_seen,last_seen,mark) VALUES(?,?,?,?,?,?)", liveID, userID, "", now, now, mark) } } } // ClearMarks resets all marks for a room. func ClearMarks(liveID string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("UPDATE users SET mark='' WHERE live_id=? AND mark!=''", liveID) } // GetMarks returns {user_id: mark} for a room. func GetMarks(liveID string) map[string]string { mu.Lock() defer mu.Unlock() c, _ := conn() rows, err := c.Query("SELECT user_id, mark FROM users WHERE live_id=? AND mark IS NOT NULL AND mark!=''", liveID) if err != nil { return nil } defer rows.Close() out := make(map[string]string) for rows.Next() { var uid, m string rows.Scan(&uid, &m) out[uid] = m } return out } // SetLead marks/unmarks a lead customer. func SetLead(liveID, userID string, lead bool) { mu.Lock() defer mu.Unlock() c, _ := conn() now := float64(time.Now().Unix()) var ll int var lts float64 if lead { ll = 1 lts = now } c.Exec("UPDATE users SET lead=?, lead_ts=? WHERE live_id=? AND user_id=?", ll, lts, liveID, userID) if lead { var id int64 c.QueryRow("SELECT id FROM users WHERE live_id=? AND user_id=?", liveID, userID).Scan(&id) if id == 0 { c.Exec("INSERT OR IGNORE INTO users(live_id,user_id,user_name,first_seen,last_seen,lead,lead_ts) VALUES(?,?,?,?,?,?,?)", liveID, userID, "", now, now, 1, now) } } } // GetLeads returns lead customers, optionally scoped to one room. func GetLeads(liveID string) []User { mu.Lock() defer mu.Unlock() c, _ := conn() var rows *sql.Rows var err error if liveID != "" { rows, err = c.Query("SELECT id,live_id,user_id,user_name,first_seen,last_seen,msg_count,gift_count,enter_count,like_count,follow_count,last_msg,last_action,hit_kw,sec_uid,avatar,mark,lead,lead_ts FROM users WHERE live_id=? AND lead=1 ORDER BY lead_ts DESC", liveID) } else { rows, err = c.Query("SELECT id,live_id,user_id,user_name,first_seen,last_seen,msg_count,gift_count,enter_count,like_count,follow_count,last_msg,last_action,hit_kw,sec_uid,avatar,mark,lead,lead_ts FROM users WHERE lead=1 ORDER BY lead_ts DESC") } if err != nil { return nil } defer rows.Close() var out []User for rows.Next() { var u User rows.Scan(&u.ID, &u.LiveID, &u.UserID, &u.UserName, &u.FirstSeen, &u.LastSeen, &u.MsgCount, &u.GiftCount, &u.EnterCount, &u.LikeCount, &u.FollowCount, &u.LastMsg, &u.LastAction, &u.HitKw, &u.SecUID, &u.Avatar, &u.Mark, &u.Lead, &u.LeadTs) out = append(out, u) } return out } var storeTypes = map[string]bool{ "chat": true, "gift": true, "enter": true, "like": true, "follow": true, "fansclub": true, "emoji_chat": true, } // StoreMsg persists a user-interaction message. func StoreMsg(liveID, msgType string, data map[string]interface{}) { if !storeTypes[msgType] { return } mu.Lock() defer mu.Unlock() c, _ := conn() ts := 0.0 if v, ok := data["timestamp"]; ok { switch x := v.(type) { case float64: ts = x case int: ts = float64(x) case int64: ts = float64(x) default: ts = float64(time.Now().Unix()) } } if ts == 0 { ts = float64(time.Now().Unix()) } payload := map[string]interface{}{"type": msgType, "room_id": liveID, "timestamp": ts} for k, v := range data { if k != "type" && k != "room_id" && k != "timestamp" { payload[k] = v } } b, _ := json.Marshal(payload) c.Exec("INSERT INTO chat_logs(live_id, ts, msg) VALUES(?,?,?)", liveID, ts, string(b)) } // GetMessages returns messages for a room with pagination. func GetMessages(liveID string, afterID int64, limit int) ([]ChatLog, int64) { mu.Lock() defer mu.Unlock() c, _ := conn() rows, err := c.Query("SELECT id, msg, seen FROM chat_logs WHERE live_id=? AND id>? ORDER BY id ASC LIMIT ?", liveID, afterID, limit) if err != nil { return nil, afterID } defer rows.Close() var out []ChatLog var lastID = afterID for rows.Next() { var cl ChatLog rows.Scan(&cl.ID, &cl.Msg, &cl.Seen) out = append(out, cl) lastID = cl.ID } return out, lastID } // SetMsgSeen marks a single message as seen. func SetMsgSeen(liveID string, msgID int64, seen bool) int { mu.Lock() defer mu.Unlock() c, _ := conn() v := 0 if seen { v = 1 } res, err := c.Exec("UPDATE chat_logs SET seen=? WHERE id=? AND live_id=?", v, msgID, liveID) if err != nil { return 0 } n, _ := res.RowsAffected() return int(n) } // SetMsgSeenMany bulk marks messages as seen. func SetMsgSeenMany(liveID string, ids []int64, seen bool) int { if len(ids) == 0 { return 0 } mu.Lock() defer mu.Unlock() c, _ := conn() if len(ids) > 2000 { ids = ids[:2000] } ph := make([]string, len(ids)) vs := make([]interface{}, len(ids)+2) vs[0] = btoi(seen) vs[1] = liveID for i, id := range ids { ph[i] = "?" vs[i+2] = id } res, err := c.Exec(fmt.Sprintf("UPDATE chat_logs SET seen=? WHERE live_id=? AND id IN (%s)", joinStrings(ph, ",")), vs...) if err != nil { return 0 } n, _ := res.RowsAffected() return int(n) } // ClearMsgSeen resets all seen flags for a room. func ClearMsgSeen(liveID string) int { mu.Lock() defer mu.Unlock() c, _ := conn() res, err := c.Exec("UPDATE chat_logs SET seen=0 WHERE live_id=? AND seen=1", liveID) if err != nil { return 0 } n, _ := res.RowsAffected() return int(n) } // GetSetting reads a setting value. func GetSetting(key string) string { mu.Lock() defer mu.Unlock() c, _ := conn() var v string c.QueryRow("SELECT v FROM settings WHERE k=?", key).Scan(&v) return v } // SetSetting writes a setting value. func SetSetting(key, val string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("INSERT OR REPLACE INTO settings(k,v) VALUES(?,?)", key, val) } // LogError stores an error log entry. func LogError(source, roomID, message, stack string) { mu.Lock() defer mu.Unlock() c, _ := conn() c.Exec("INSERT INTO error_logs(ts, source, room_id, message, stack) VALUES(?,?,?,?,?)", float64(time.Now().Unix()), truncStr(source, 20), truncStr(roomID, 40), truncStr(message, 2000), truncStr(stack, 8000), ) } // GetErrors returns recent error logs. func GetErrors(limit int) []ErrorLog { mu.Lock() defer mu.Unlock() c, _ := conn() rows, err := c.Query("SELECT id, ts, source, room_id, message, stack FROM error_logs ORDER BY id DESC LIMIT ?", limit) if err != nil { return nil } defer rows.Close() var out []ErrorLog for rows.Next() { var el ErrorLog rows.Scan(&el.ID, &el.TS, &el.Source, &el.RoomID, &el.Message, &el.Stack) out = append(out, el) } return out } func truncStr(s string, n int) string { if len(s) <= n { return s } return s[:n] } func btoi(b bool) int { if b { return 1 } return 0 }