package fetcher import ( "compress/gzip" "context" "encoding/binary" "encoding/json" "fmt" "io" "log" "math/rand" "net/http" "net/url" "regexp" "strings" "sync" "sync/atomic" "time" "bytes" "nhooyr.io/websocket" "github.com/qingwuyingxing/douyin-live-fetcher/internal/broadcaster" "github.com/qingwuyingxing/douyin-live-fetcher/internal/db" "github.com/qingwuyingxing/douyin-live-fetcher/internal/sign" ) type Fetcher struct { LiveID string bc *broadcaster.Broadcaster ok atomic.Bool ws atomic.Value lastPkt int64 silent bool kws string kr *regexp.Regexp users map[string]map[string]interface{} lk sync.RWMutex done chan struct{} } var didC int64 func NewFetcher(liveID string, bc *broadcaster.Broadcaster, keywords string) *Fetcher { var kr *regexp.Regexp if keywords != "" { parts := strings.FieldsFunc(keywords, func(c rune) bool { return c == ',' || c == ' ' || c == '\n' }) if len(parts) > 0 { escaped := make([]string, len(parts)) for i, p := range parts { escaped[i] = regexp.QuoteMeta(p) } kr = regexp.MustCompile("(?i)" + strings.Join(escaped, "|")) } } f := &Fetcher{ LiveID: liveID, bc: bc, lastPkt: time.Now().Unix(), kws: keywords, kr: kr, users: make(map[string]map[string]interface{}), done: make(chan struct{}), } for _, u := range db.GetUsers(liveID) { f.users[u.UserID] = map[string]interface{}{ "user_id": u.UserID, "user_name": u.UserName, "msg_count": u.MsgCount, "gift_count": u.GiftCount, "enter_count": u.EnterCount, "like_count": u.LikeCount, "follow_count": u.FollowCount, "last_msg": u.LastMsg, "last_action": u.LastAction, "hit_kw": u.HitKw, "sec_uid": u.SecUID, "avatar": u.Avatar, "mark": u.Mark, "lead": u.Lead, } } return f } func (f *Fetcher) Start() { go f._connectWebSocket() go f._watchdog() } func (f *Fetcher) Stop() { select { case <-f.done: return default: close(f.done) } if ws, ok := f.ws.Load().(*websocket.Conn); ok && ws != nil { cancel := func() {} _ = cancel ws.CloseNow() } f.ok.Store(false) } func (f *Fetcher) Ok() bool { return f.ok.Load() } func (f *Fetcher) LastPkt() float64 { return float64(atomic.LoadInt64(&f.lastPkt)) } func (f *Fetcher) IsSilent() bool { return f.silent } func (f *Fetcher) Keywords() string { return f.kws } func (f *Fetcher) SetKeywords(kws string) { f.kws = kws var kr *regexp.Regexp if kws != "" { parts := strings.FieldsFunc(kws, func(c rune) bool { return c == ',' || c == ' ' || c == '\n' }) if len(parts) > 0 { escaped := make([]string, len(parts)) for i, p := range parts { escaped[i] = regexp.QuoteMeta(p) } kr = regexp.MustCompile("(?i)" + strings.Join(escaped, "|")) } } f.kr = kr } func (f *Fetcher) SetMark(uid, mark string) { db.SetMark(f.LiveID, uid, mark) } func (f *Fetcher) ClearMarks() { db.ClearMarks(f.LiveID) } func (f *Fetcher) SetLead(uid string, lead bool) { db.SetLead(f.LiveID, uid, lead) } func (f *Fetcher) ThreadAlive() bool { select { case <-f.done: return false default: return true } } func (f *Fetcher) Users() []map[string]interface{} { f.lk.RLock() defer f.lk.RUnlock() out := make([]map[string]interface{}, 0, len(f.users)) for _, u := range f.users { out = append(out, u) } return out } func (f *Fetcher) _connectWebSocket() { for { select { case <-f.done: return default: } ttwid := f.getTTWID() roomID := f.getRoomID(ttwid) if roomID == "" { roomID = f.LiveID } nonce := f.getACNonce() sig := sign.GetACSignature("www.douyin.com", nonce, uaStr, 0) wssURL := f.buildWSSURL(roomID, sig, nonce, ttwid) sig2, err := sign.GenerateSignature(wssURL) if err != nil { log.Printf("[fetcher:%s] sign error: %v", f.LiveID, err) } wssURL += "&signature=" + sig2 f._wsLoop(wssURL) select { case <-f.done: return case <-time.After(5 * time.Second): } } } const uaStr = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36" func (f *Fetcher) getTTWID() string { req, _ := http.NewRequest("GET", "https://live.douyin.com/"+f.LiveID, nil) req.Header.Set("User-Agent", uaStr) client := &http.Client{Timeout: 5 * time.Second, CheckRedirect: func(req *http.Request, via []*http.Request) error { return http.ErrUseLastResponse }} resp, err := client.Do(req) if err != nil { return "" } defer resp.Body.Close() for _, c := range resp.Cookies() { if c.Name == "ttwid" { return c.Value } } return "" } func (f *Fetcher) getRoomID(ttwid string) string { req, _ := http.NewRequest("GET", "https://live.douyin.com/"+f.LiveID, nil) req.Header.Set("User-Agent", uaStr) req.Header.Set("Cookie", "ttwid="+ttwid+"; __ac_nonce=0123407cc00a9e438deb4") client := &http.Client{Timeout: 5 * time.Second} resp, err := client.Do(req) if err != nil { return "" } defer resp.Body.Close() body, _ := io.ReadAll(resp.Body) m := regexp.MustCompile(`roomId[":\s]+(\d+)`).FindSubmatch(body) if len(m) > 1 { return string(m[1]) } m = regexp.MustCompile(`"room_id_str":"(\d+)"`).FindSubmatch(body) if len(m) > 1 { return string(m[1]) } return "" } func (f *Fetcher) getACNonce() string { client := &http.Client{Timeout: 5 * time.Second} resp, err := client.Get("https://www.douyin.com/") if err != nil { return "" } defer resp.Body.Close() for _, c := range resp.Cookies() { if c.Name == "__ac_nonce" { return c.Value } } return "" } func (f *Fetcher) buildWSSURL(roomID, sig, nonce, ttwid string) string { did := fmt.Sprintf("7%d", atomic.AddInt64(&didC, 1)) nowMs := time.Now().UnixMilli() internalExt := fmt.Sprintf("internal_src:dim|wss_push_room_id:%s|wss_push_did:%s|first_req_ms:%d|fetch_time:%d|seq:1|wss_info:0-%d-0-0|wrds_v:%d", roomID, did, nowMs, nowMs, nowMs, rand.Int63n(9999999999999)) params := url.Values{} params.Set("app_name", "douyin_web") params.Set("version_code", "180800") params.Set("webcast_sdk_version", "1.0.14-beta.0") params.Set("compress", "gzip") params.Set("device_platform", "web") params.Set("cookie_enabled", "true") params.Set("screen_width", "1536") params.Set("screen_height", "864") params.Set("browser_language", "zh-CN") params.Set("browser_platform", "Win32") params.Set("browser_name", "Mozilla") params.Set("browser_version", "5.0%20(Windows%20NT%2010.0;%20Win64;%20x64)%20AppleWebKit/537.36%20(KHTML,%20like%20Gecko)%20Chrome/126.0.0.0%20Safari/537.36") params.Set("browser_online", "true") params.Set("tz_name", "Asia/Shanghai") params.Set("cursor", "d-1_u-1_fh-1_t-1_r-1") params.Set("internal_ext", internalExt) params.Set("host", "https://live.douyin.com") params.Set("aid", "6383") params.Set("live_id", "1") params.Set("did_rule", "3") params.Set("endpoint", "live_pc") params.Set("support_wrds", "1") params.Set("user_unique_id", did) params.Set("im_path", "/webcast/im/fetch/") params.Set("identity", "audience") params.Set("need_persist_msg_count", "15") params.Set("room_id", roomID) params.Set("heartbeatDuration", "0") base := "wss://webcast100-ws-web-lq.douyin.com/webcast/im/push/v2/" return base + "?" + params.Encode() } func (f *Fetcher) _wsLoop(wssURL string) { headers := make(http.Header) headers.Set("User-Agent", uaStr) headers.Set("Referer", "https://live.douyin.com/"+f.LiveID) ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) ws, _, err := websocket.Dial(ctx, wssURL, &websocket.DialOptions{HTTPHeader: headers}) cancel() if err != nil { log.Printf("[fetcher:%s] ws dial error: %v", f.LiveID, err) return } f.ws.Store(ws) f.ok.Store(true) atomic.StoreInt64(&f.lastPkt, time.Now().Unix()) f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "connected"}) go f._heartbeat(ws) go f._readLoop(ws) } func (f *Fetcher) _heartbeat(ws *websocket.Conn) { tick := time.NewTicker(5 * time.Second) defer tick.Stop() for { select { case <-f.done: return case <-tick.C: pushFrame := buildPushFrame(1, []byte{}) if err := ws.Write(context.Background(), websocket.MessageBinary, pushFrame); err != nil { return } f._watchdogCheck() } } } func buildPushFrame(msgType byte, payload []byte) []byte { length := uint32(16 + len(payload)) buf := make([]byte, length) binary.BigEndian.PutUint32(buf[0:4], 2) buf[4] = msgType binary.BigEndian.PutUint32(buf[8:12], length) binary.BigEndian.PutUint64(buf[12:16], uint64(time.Now().UnixNano())) copy(buf[16:], payload) return buf } func (f *Fetcher) _readLoop(ws *websocket.Conn) { defer func() { f.ok.Store(false) f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "disconnected"}) ws.Close(websocket.StatusNormalClosure, "") }() buf := make([]byte, 4*1024*1024) for { select { case <-f.done: return default: } msgType, reader, err := ws.Reader(contextWithTimeout(30 * time.Second)) if err != nil { return } if msgType != websocket.MessageBinary { continue } n, err := io.ReadFull(reader, buf) if err != nil { return } data := buf[:n] if len(data) < 16 { continue } atomic.StoreInt64(&f.lastPkt, time.Now().Unix()) f._watchdogCheck() if err := f._dispatch(data); err != nil { log.Printf("[fetcher:%s] dispatch error: %v", f.LiveID, err) } } } func (f *Fetcher) _dispatch(data []byte) error { if len(data) < 16 { return nil } version := binary.BigEndian.Uint32(data[0:4]) msgTypeByte := data[4] compressed := (msgTypeByte & 0x80) != 0 needsAck := (msgTypeByte & 0x40) != 0 msgType := msgTypeByte & 0x3F length := binary.BigEndian.Uint32(data[8:12]) logID := binary.BigEndian.Uint64(data[12:16]) var payload []byte offset := 16 if compressed { switch msgType { case 0: decomp, err := decompressGzip(data[offset:]) if err != nil { return err } payload = decomp case 1: payload = data[offset:] default: return nil } } else { payload = data[offset:] } if needsAck && length > 0 { ack := make([]byte, 16+len(payload)) binary.BigEndian.PutUint32(ack[0:4], version) ack[4] = msgTypeByte | 0x40 binary.BigEndian.PutUint32(ack[8:12], uint32(len(payload))) binary.BigEndian.PutUint64(ack[12:16], logID) copy(ack[16:], payload) if ws, ok := f.ws.Load().(*websocket.Conn); ok && ws != nil { ws.Write(context.Background(), websocket.MessageBinary, ack) } } var msgs []dyMessage if err := json.Unmarshal(payload, &msgs); err != nil { return nil } for _, msg := range msgs { switch msg.Method { case "WebcastChatMessage": f._parseChatMsg(msg.Payload) case "WebcastGiftMessage": f._parseGiftMsg(msg.Payload) case "WebcastMemberMessage": f._parseMemberMsg(msg.Payload) case "WebcastLikeMessage": f._parseLikeMsg(msg.Payload) case "WebcastSocialMessage": f._parseSocialMsg(msg.Payload) case "WebcastFansclubMessage": f._parseFansclubMsg(msg.Payload) case "WebcastControlMessage": f._parseControlMsg(msg.Payload) case "WebcastEmojiChatMessage": f._parseEmojiChatMsg(msg.Payload) case "WebcastRoomUserSeqMessage": f._parseRoomUserSeqMsg(msg.Payload) } } return nil } type dyMessage struct { Method string `json:"method"` Payload json.RawMessage `json:"payload"` } func decompressGzip(data []byte) ([]byte, error) { r, err := gzip.NewReader(bytes.NewReader(data)) if err != nil { return nil, err } defer r.Close() return io.ReadAll(r) } func (f *Fetcher) _parseChatMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` Content string `json:"content"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec := m.User.SecUID av := "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } kw := "" if f.kr != nil && f.kr.MatchString(m.Content) { kw = f.kr.FindString(m.Content) } f.bc.Broadcast(map[string]interface{}{"type": "chat", "room_id": f.LiveID, "user_name": m.User.NickName, "user_id": uid, "sec_uid": sec, "avatar": av, "content": m.Content, "keyword": kw}) db.StoreMsg(f.LiveID, "chat", map[string]interface{}{"type": "chat", "timestamp": float64(time.Now().Unix()), "user_name": m.User.NickName, "user_id": uid, "content": m.Content, "keyword": kw}) f._upsertUser(uid, m.User.NickName, "chat", map[string]string{"content": m.Content, "keyword": kw, "sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseGiftMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` Gift struct{ Name string `json:"name"` } `json:"gift"` ComboCount int `json:"comboCount"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec, av := m.User.SecUID, "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } f.bc.Broadcast(map[string]interface{}{"type": "gift", "room_id": f.LiveID, "user_name": m.User.NickName, "user_id": uid, "sec_uid": sec, "avatar": av, "gift_name": m.Gift.Name, "count": m.ComboCount}) db.StoreMsg(f.LiveID, "gift", map[string]interface{}{"type": "gift", "timestamp": float64(time.Now().Unix()), "user_name": m.User.NickName, "user_id": uid, "gift_name": m.Gift.Name, "count": m.ComboCount}) f._upsertUser(uid, m.User.NickName, "gift", map[string]string{"sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseMemberMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` Gender int `json:"gender"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec, av := m.User.SecUID, "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } gender := "?" if m.User.Gender == 0 { gender = "女" } else if m.User.Gender == 1 { gender = "男" } f.bc.Broadcast(map[string]interface{}{"type": "enter", "room_id": f.LiveID, "user_name": m.User.NickName, "user_id": uid, "gender": gender, "sec_uid": sec, "avatar": av}) db.StoreMsg(f.LiveID, "enter", map[string]interface{}{"type": "enter", "timestamp": float64(time.Now().Unix()), "user_name": m.User.NickName, "user_id": uid, "gender": gender}) f._upsertUser(uid, m.User.NickName, "enter", map[string]string{"sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseLikeMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` Count int `json:"count"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec, av := m.User.SecUID, "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } f.bc.Broadcast(map[string]interface{}{"type": "like", "room_id": f.LiveID, "user_name": m.User.NickName, "user_id": uid, "sec_uid": sec, "avatar": av, "count": m.Count}) db.StoreMsg(f.LiveID, "like", map[string]interface{}{"type": "like", "timestamp": float64(time.Now().Unix()), "user_name": m.User.NickName, "user_id": uid, "count": m.Count}) f._upsertUser(uid, m.User.NickName, "like", map[string]string{"sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseSocialMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec, av := m.User.SecUID, "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } f.bc.Broadcast(map[string]interface{}{"type": "follow", "room_id": f.LiveID, "user_name": m.User.NickName, "user_id": uid, "sec_uid": sec, "avatar": av}) db.StoreMsg(f.LiveID, "follow", map[string]interface{}{"type": "follow", "timestamp": float64(time.Now().Unix()), "user_name": m.User.NickName, "user_id": uid}) f._upsertUser(uid, m.User.NickName, "follow", map[string]string{"sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseFansclubMsg(payload []byte) { var m struct { User struct { NickName string `json:"nickName"` ID uint64 `json:"id"` SecUID string `json:"secUid"` AvatarThumb struct { URLList []string `json:"urlList"` } `json:"avatarThumb"` } `json:"user"` Content string `json:"content"` Type int `json:"type"` } json.Unmarshal(payload, &m) uid := fmt.Sprintf("%d", m.User.ID) sec, av := m.User.SecUID, "" if len(m.User.AvatarThumb.URLList) > 0 { av = m.User.AvatarThumb.URLList[0] } f.bc.Broadcast(map[string]interface{}{"type": "fansclub", "room_id": f.LiveID, "content": m.Content, "user_name": m.User.NickName, "user_id": uid, "sec_uid": sec, "avatar": av, "fc_type": m.Type}) f._upsertUser(uid, m.User.NickName, "follow", map[string]string{"sec_uid": sec, "avatar": av}) } func (f *Fetcher) _parseControlMsg(payload []byte) { var m struct { Status int `json:"status"` } json.Unmarshal(payload, &m) if m.Status == 3 { f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "stream_ended"}) go f.Stop() } } func (f *Fetcher) _parseEmojiChatMsg(payload []byte) { var m struct { EmojiID int `json:"emojiId"` DefaultContent string `json:"defaultContent"` User struct{ NickName string `json:"nickName"` } `json:"user"` } json.Unmarshal(payload, &m) f.bc.Broadcast(map[string]interface{}{"type": "emoji_chat", "room_id": f.LiveID, "emoji_id": m.EmojiID, "user_name": m.User.NickName, "default_content": m.DefaultContent}) } func (f *Fetcher) _parseRoomUserSeqMsg(payload []byte) { var m struct { Total int `json:"total"` TotalPvForAnchor int `json:"totalPvForAnchor"` } json.Unmarshal(payload, &m) f.bc.Broadcast(map[string]interface{}{"type": "stats", "room_id": f.LiveID, "current_viewers": m.Total, "total_viewers": m.TotalPvForAnchor}) } func (f *Fetcher) _upsertUser(uid, uname, action string, extra map[string]string) { db.UpsertUser(f.LiveID, uid, uname, action, extra) f.lk.Lock() defer f.lk.Unlock() if _, ok := f.users[uid]; !ok { f.users[uid] = make(map[string]interface{}) } u := f.users[uid] u["user_id"] = uid u["user_name"] = uname u["last_action"] = action u["last_seen"] = float64(time.Now().Unix()) if v, ok := extra["content"]; ok && v != "" { u["msg_count"] = asInt(u["msg_count"]) + 1 u["last_msg"] = v } if v, ok := extra["keyword"]; ok && v != "" { u["hit_kw"] = v } if v, ok := extra["sec_uid"]; ok && v != "" { u["sec_uid"] = v } if v, ok := extra["avatar"]; ok && v != "" { u["avatar"] = v } switch action { case "gift": u["gift_count"] = asInt(u["gift_count"]) + 1 case "enter": u["enter_count"] = asInt(u["enter_count"]) + 1 case "like": u["like_count"] = asInt(u["like_count"]) + 1 case "follow": u["follow_count"] = asInt(u["follow_count"]) + 1 } f.users[uid] = u } func asInt(v interface{}) int { switch x := v.(type) { case int: return x case int64: return int(x) case float64: return int(x) case nil: return 0 } return 0 } func (f *Fetcher) _watchdog() { tick := time.NewTicker(1 * time.Second) stTick := time.NewTicker(25 * time.Second) defer tick.Stop() defer stTick.Stop() for { select { case <-f.done: return case <-tick.C: f._watchdogCheck() case <-stTick.C: if f.ok.Load() && !f.silent { f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "connected"}) } } } } func (f *Fetcher) _watchdogCheck() { now := time.Now().Unix() dt := now - atomic.LoadInt64(&f.lastPkt) if !f.ok.Load() { return } if dt > 60 && !f.silent { f.silent = true f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "stream_ended"}) log.Printf("[fetcher:%s] no data for %.0fs, stream_ended", f.LiveID, float64(dt)) } else if dt <= 40 && f.silent { f.silent = false f.bc.Broadcast(map[string]interface{}{"type": "status", "room_id": f.LiveID, "content": "connected"}) log.Printf("[fetcher:%s] data resumed, connected", f.LiveID) } } func contextWithTimeout(d time.Duration) context.Context { ctx, cancel := context.WithTimeout(context.Background(), d) go func() { <-ctx.Done() cancel() }() return ctx }