redis.go 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813
  1. package dao
  2. import (
  3. "context"
  4. "fmt"
  5. "go-common/app/service/video/stream-mng/model"
  6. "go-common/library/cache/redis"
  7. "go-common/library/log"
  8. "strconv"
  9. "strings"
  10. )
  11. /*
  12. kv结构:
  13. streamName: room_id
  14. hash结构如下:
  15. // 存储所有的流名
  16. room_id + name + name1,name2, name3
  17. // 默认直推
  18. room_id + name1:default + src
  19. // 一个流名下的直推
  20. room_id + name1:origin + src
  21. // 一个流名下的转推
  22. room_id + name1:streaming + src
  23. // 一个流名下的key
  24. room_id + name1:key + key
  25. 。。。
  26. */
  27. const (
  28. // kv => name: rid
  29. _streamNameKey = "mng:name:%s"
  30. // hash key => rid: field :value
  31. _streamRooIDKey = "mng:rid:%d"
  32. _streamRoomFieldAllName = "mng:name:all"
  33. _streamRoomFieldDefault = "%s:default"
  34. _streamRoomFieldOrigin = "%s:origin"
  35. _streamRoomFieldForward = "%s:forward"
  36. _streamRoomFieldSecret = "%s:key"
  37. _streamRoomFieldOption = "%s:options"
  38. _streamRoomFieldHot = "%d:hot"
  39. _streamLastCDN = "last:cdn:%d"
  40. // 切流记录
  41. _streamChangeSrc = "change:src:%d"
  42. //房间冷热流
  43. _streamRoomHot = "room:hot:%d"
  44. // 存12小时
  45. _streamExpireTime = 4 * 3600
  46. // stream_name map room_id 流名对应房间号 这个是不变的
  47. _nameExpireTime = 365 * 86400
  48. // 存一年
  49. _lastCDNExpireTime = 365 * 86400
  50. _changeSrcExpireTime = 365 * 86400
  51. )
  52. func (d *Dao) getStreamNamekey(streamName string) string {
  53. return fmt.Sprintf(_streamNameKey, streamName)
  54. }
  55. func (d *Dao) getRoomIDKey(rid int64) string {
  56. return fmt.Sprintf(_streamRooIDKey, rid)
  57. }
  58. func (d *Dao) getRoomFieldDefaultKey(streamName string) string {
  59. return fmt.Sprintf(_streamRoomFieldDefault, streamName)
  60. }
  61. func (d *Dao) getRoomFieldOriginKey(streamName string) string {
  62. return fmt.Sprintf(_streamRoomFieldOrigin, streamName)
  63. }
  64. func (d *Dao) getRoomFieldForwardKey(streamName string) string {
  65. return fmt.Sprintf(_streamRoomFieldForward, streamName)
  66. }
  67. func (d *Dao) getRoomFieldSecretKey(streamName string) string {
  68. return fmt.Sprintf(_streamRoomFieldSecret, streamName)
  69. }
  70. func (d *Dao) getRoomFieldOption(streamName string) string {
  71. return fmt.Sprintf(_streamRoomFieldOption, streamName)
  72. }
  73. func (d *Dao) getRoomFieldHotKey(rid int64) string {
  74. return fmt.Sprintf(_streamRoomFieldHot, rid)
  75. }
  76. func (d *Dao) getRoomHotKey(rid int64) string {
  77. return fmt.Sprintf(_streamRoomHot, rid)
  78. }
  79. func (d *Dao) getLastCDNKey(rid int64) string {
  80. return fmt.Sprintf(_streamLastCDN, rid)
  81. }
  82. func (d *Dao) getChangeSrcKey(rid int64) string {
  83. return fmt.Sprintf(_streamChangeSrc, rid)
  84. }
  85. // CacheStreamFullInfo 从缓存取流信息, 可传入流名, 也可以传入rid
  86. func (d *Dao) CacheStreamFullInfo(c context.Context, rid int64, sname string) (res *model.StreamFullInfo, err error) {
  87. if sname != "" {
  88. infos, err := d.CacheStreamRIDByName(c, sname)
  89. if err != nil {
  90. return nil, err
  91. }
  92. if infos == nil || infos.RoomID <= 0 {
  93. return nil, fmt.Errorf("can not find any info by sname =%s", sname)
  94. }
  95. rid = infos.RoomID
  96. }
  97. // 先从本地缓存中取
  98. res = d.loadStreamInfo(c, rid)
  99. if res != nil {
  100. log.Warn("get from local cache")
  101. return res, nil
  102. }
  103. conn := d.redis.Get(c)
  104. defer conn.Close()
  105. roomKey := d.getRoomIDKey(rid)
  106. hotKey := d.getRoomFieldHotKey(rid)
  107. // 先判断过期时间是否为-1, 为-1则删除
  108. //ttl, _ := redis.Int(conn.Do("TTL", roomKey))
  109. //if ttl == -1 {
  110. // d.DeleteStreamByRIDFromCache(c, rid)
  111. // return nil, nil
  112. //}
  113. values, err := redis.StringMap(conn.Do("HGETALL", roomKey))
  114. log.Warn("%v", values)
  115. if err != nil {
  116. return nil, err
  117. }
  118. if len(values) == 0 {
  119. return nil, nil
  120. }
  121. resp := model.StreamFullInfo{}
  122. hot, _ := strconv.ParseInt(values[hotKey], 10, 64)
  123. resp.Hot = hot
  124. allNames := strings.Split(values[_streamRoomFieldAllName], "|")
  125. for _, n := range allNames {
  126. defaultUpStream := d.getRoomFieldDefaultKey(n)
  127. originKey := d.getRoomFieldOriginKey(n)
  128. forwardKey := d.getRoomFieldForwardKey(n)
  129. key := d.getRoomFieldSecretKey(n)
  130. options := d.getRoomFieldOption(n)
  131. if values[originKey] != "" && values[forwardKey] != "" && values[key] != "" && values[defaultUpStream] != "" {
  132. resp.RoomID = rid
  133. base := model.StreamBase{}
  134. base.StreamName = n
  135. or, _ := strconv.ParseInt(values[originKey], 10, 64)
  136. base.Origin = or
  137. de, _ := strconv.ParseInt(values[defaultUpStream], 10, 64)
  138. base.DefaultUpStream = de
  139. if de == 0 {
  140. return nil, fmt.Errorf("default is 0")
  141. }
  142. forward, _ := strconv.ParseInt(values[forwardKey], 10, 64)
  143. var num int64
  144. for num = 256; num > 0; num /= 2 {
  145. if ((forward & num) == num) && (num != or) {
  146. base.Forward = append(base.Forward, num)
  147. }
  148. }
  149. op, _ := strconv.ParseInt(values[options], 10, 64)
  150. base.Options = op
  151. //这里判断是否有wmask mmask的流
  152. if 4&op == 4 {
  153. base.Wmask = true
  154. }
  155. if 8&op == 8 {
  156. base.Mmask = true
  157. }
  158. base.Key = values[key]
  159. if strings.Contains(n, "_bs_") {
  160. base.Type = 2
  161. } else {
  162. base.Type = 1
  163. }
  164. resp.List = append(resp.List, &base)
  165. }
  166. }
  167. if len(resp.List) > 0 {
  168. // 存储到local cache
  169. d.storeStreamInfo(c, &resp)
  170. return &resp, nil
  171. }
  172. return nil, nil
  173. }
  174. // AddCacheStreamFullInfo 修改缓存数据
  175. func (d *Dao) AddCacheStreamFullInfo(c context.Context, id int64, stream *model.StreamFullInfo) error {
  176. if stream == nil || stream.RoomID <= 0 {
  177. return nil
  178. }
  179. conn := d.redis.Get(c)
  180. defer func() {
  181. //conn.Do("EXPIRE", d.getRoomIDKey(id), _streamExpireTime)
  182. conn.Close()
  183. }()
  184. streamExpireTime := _streamExpireTime
  185. rid := stream.RoomID
  186. roomKey := d.getRoomIDKey(rid)
  187. allName := ""
  188. len := 0
  189. for _, v := range stream.List {
  190. len++
  191. allName = fmt.Sprintf("%s%s|", allName, v.StreamName)
  192. // kv 设置流名和room_id映射关系
  193. nameKey := d.getStreamNamekey(v.StreamName)
  194. if err := conn.Send("SET", nameKey, rid); err != nil {
  195. return fmt.Errorf("conn.Do(set, %s, %d) error(%v)", nameKey, rid, err)
  196. }
  197. if err := conn.Send("EXPIRE", nameKey, _nameExpireTime); err != nil {
  198. return fmt.Errorf("conn.Do(EXPIRE, %s,%d) error(%v)", nameKey, _nameExpireTime, err)
  199. }
  200. // hash 设置room_id下的field 和key
  201. field := d.getRoomFieldDefaultKey(v.StreamName)
  202. if v.DefaultUpStream == 0 {
  203. return fmt.Errorf("rid= %v, default is 0", roomKey)
  204. }
  205. if err := conn.Send("HSET", roomKey, field, v.DefaultUpStream); err != nil {
  206. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.DefaultUpStream, err)
  207. }
  208. if len == 1 {
  209. if err := conn.Send("EXPIRE", roomKey, streamExpireTime); err != nil {
  210. log.Infov(c, log.KV("conn.EXPIRE error", err.Error()))
  211. return fmt.Errorf("conn.Do(EXPIRE, %s, %d) error(%v)", roomKey, streamExpireTime, err)
  212. }
  213. }
  214. field = d.getRoomFieldOriginKey(v.StreamName)
  215. if err := conn.Send("HSET", roomKey, field, v.Origin); err != nil {
  216. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Origin, err)
  217. }
  218. field = d.getRoomFieldForwardKey(v.StreamName)
  219. var num int64
  220. for _, f := range v.Forward {
  221. num += f
  222. }
  223. if err := conn.Send("HSET", roomKey, field, num); err != nil {
  224. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Forward, err)
  225. }
  226. field = d.getRoomFieldSecretKey(v.StreamName)
  227. if err := conn.Send("HSET", roomKey, field, v.Key); err != nil {
  228. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Key, err)
  229. }
  230. field = d.getRoomFieldOption(v.StreamName)
  231. if err := conn.Send("HSET", roomKey, field, v.Options); err != nil {
  232. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Options, err)
  233. }
  234. }
  235. // 去除最后的|
  236. allName = strings.Trim(allName, "|")
  237. //log.Warn("%v", allName)
  238. if err := conn.Send("HSET", roomKey, _streamRoomFieldAllName, allName); err != nil {
  239. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, _streamRoomFieldAllName, allName, err)
  240. }
  241. if err := conn.Flush(); err != nil {
  242. log.Infov(c, log.KV("conn.Flush error(%v)", err.Error()))
  243. return fmt.Errorf("conn.Flush error(%v)", err)
  244. }
  245. for i := 0; i < 7*len+2; i++ {
  246. if _, err := conn.Receive(); err != nil {
  247. log.Infov(c, log.KV("conn.Receive error(%v)", err.Error()))
  248. return fmt.Errorf("conn.Receive error(%v)", err)
  249. }
  250. }
  251. return nil
  252. }
  253. // CacheStreamRIDByName 根据流名查房间号
  254. func (d *Dao) CacheStreamRIDByName(c context.Context, sname string) (res *model.StreamFullInfo, err error) {
  255. conn := d.redis.Get(c)
  256. defer conn.Close()
  257. nameKey := d.getStreamNamekey(sname)
  258. rid, err := redis.Int64(conn.Do("GET", nameKey))
  259. if err != nil {
  260. return nil, err
  261. }
  262. if rid <= 0 {
  263. return nil, nil
  264. }
  265. res = &model.StreamFullInfo{
  266. RoomID: rid,
  267. }
  268. return res, nil
  269. }
  270. // AddCacheStreamRIDByName 增加缓存
  271. func (d *Dao) AddCacheStreamRIDByName(c context.Context, sname string, stream *model.StreamFullInfo) error {
  272. return d.AddCacheStreamFullInfo(c, 0, stream)
  273. }
  274. // AddCacheMultiStreamInfo 批量增加redis
  275. func (d *Dao) AddCacheMultiStreamInfo(c context.Context, res map[int64]*model.StreamFullInfo) error {
  276. conn := d.redis.Get(c)
  277. defer func() {
  278. //for _, stream := range res {
  279. // conn.Do("EXPIRE", d.getRoomIDKey(stream.RoomID), _streamExpireTime)
  280. //}
  281. conn.Close()
  282. }()
  283. count := 0
  284. for _, stream := range res {
  285. streamExpireTime := _streamExpireTime
  286. rid := stream.RoomID
  287. roomKey := d.getRoomIDKey(rid)
  288. allName := ""
  289. len := 0
  290. for _, v := range stream.List {
  291. len++
  292. allName = fmt.Sprintf("%s%s|", allName, v.StreamName)
  293. // kv 设置流名和room_id映射关系
  294. nameKey := d.getStreamNamekey(v.StreamName)
  295. if err := conn.Send("SET", nameKey, rid); err != nil {
  296. return fmt.Errorf("conn.Do(set, %s, %d) error(%v)", nameKey, rid, err)
  297. }
  298. if err := conn.Send("EXPIRE", nameKey, _nameExpireTime); err != nil {
  299. return fmt.Errorf("conn.Do(EXPIRE, %s,%d) error(%v)", nameKey, _nameExpireTime, err)
  300. }
  301. // hash 设置room_id下的field 和key
  302. field := d.getRoomFieldDefaultKey(v.StreamName)
  303. if err := conn.Send("HSET", roomKey, field, v.DefaultUpStream); err != nil {
  304. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.DefaultUpStream, err)
  305. }
  306. if len == 1 {
  307. if err := conn.Send("EXPIRE", roomKey, streamExpireTime); err != nil {
  308. return fmt.Errorf("conn.Do(EXPIRE, %s, %d) error(%v)", roomKey, streamExpireTime, err)
  309. }
  310. }
  311. field = d.getRoomFieldOriginKey(v.StreamName)
  312. if err := conn.Send("HSET", roomKey, field, v.Origin); err != nil {
  313. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Origin, err)
  314. }
  315. field = d.getRoomFieldForwardKey(v.StreamName)
  316. var num int64
  317. for _, f := range v.Forward {
  318. num += f
  319. }
  320. if err := conn.Send("HSET", roomKey, field, num); err != nil {
  321. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Forward, err)
  322. }
  323. field = d.getRoomFieldSecretKey(v.StreamName)
  324. if err := conn.Send("HSET", roomKey, field, v.Key); err != nil {
  325. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Key, err)
  326. }
  327. field = d.getRoomFieldOption(v.StreamName)
  328. if err := conn.Send("HSET", roomKey, field, v.Options); err != nil {
  329. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, field, v.Options, err)
  330. }
  331. }
  332. // 去除最后的|
  333. allName = strings.Trim(allName, "|")
  334. if err := conn.Send("HSET", roomKey, _streamRoomFieldAllName, allName); err != nil {
  335. return fmt.Errorf("conn.Do(HSET, %s, %s, %v) error(%v)", roomKey, _streamRoomFieldAllName, allName, err)
  336. }
  337. count += (7*len + 2)
  338. }
  339. if err := conn.Flush(); err != nil {
  340. return fmt.Errorf("conn.Flush error(%v)", err)
  341. }
  342. // 这里要len
  343. for i := 0; i < count; i++ {
  344. if _, err := conn.Receive(); err != nil {
  345. return fmt.Errorf("conn.Receive error(%v)", err)
  346. }
  347. }
  348. return nil
  349. }
  350. // CacheMultiStreamInfo 批量从redis中获取数据
  351. func (d *Dao) CacheMultiStreamInfo(c context.Context, rids []int64) (res map[int64]*model.StreamFullInfo, err error) {
  352. if len(rids) == 0 {
  353. return nil, nil
  354. }
  355. infos := map[int64]*model.StreamFullInfo{}
  356. // 先从local cache读取
  357. localInfos, missRids := d.loadMultiStreamInfo(c, rids)
  358. //log.Warn("%v=%v", localInfos, missRids)
  359. // 若全部命中local cache,直接返回
  360. if len(missRids) == 0 {
  361. log.Warn("all hit local cache")
  362. return localInfos, nil
  363. }
  364. if len(localInfos) != 0 {
  365. infos = localInfos
  366. }
  367. conn := d.redis.Get(c)
  368. defer conn.Close()
  369. for _, id := range missRids {
  370. key := d.getRoomIDKey(id)
  371. if err := conn.Send("HGETALL", key); err != nil {
  372. log.Errorv(c, log.KV("log", fmt.Sprintf("redis: conn.Send(HGETALL, %s) error(%v)", key, err)))
  373. return nil, fmt.Errorf("redis: conn.Send(HGETALL, %s) error(%v)", key, err)
  374. }
  375. }
  376. if err := conn.Flush(); err != nil {
  377. return nil, fmt.Errorf("redis: conn.Flush error(%v)", err)
  378. }
  379. for i := 0; i < len(missRids); i++ {
  380. if values, err := redis.StringMap(conn.Receive()); err == nil {
  381. if len(values) == 0 {
  382. continue
  383. }
  384. item := model.StreamFullInfo{}
  385. hotKey := d.getRoomFieldHotKey(missRids[i])
  386. hot, _ := strconv.ParseInt(values[hotKey], 10, 64)
  387. item.Hot = hot
  388. allNames := strings.Split(values[_streamRoomFieldAllName], "|")
  389. for _, n := range allNames {
  390. defaultUpStream := d.getRoomFieldDefaultKey(n)
  391. originKey := d.getRoomFieldOriginKey(n)
  392. forwardKey := d.getRoomFieldForwardKey(n)
  393. key := d.getRoomFieldSecretKey(n)
  394. options := d.getRoomFieldOption(n)
  395. if values[originKey] != "" && values[forwardKey] != "" && values[key] != "" {
  396. item.RoomID = missRids[i]
  397. base := model.StreamBase{}
  398. base.StreamName = n
  399. or, _ := strconv.ParseInt(values[originKey], 10, 64)
  400. base.Origin = or
  401. de, _ := strconv.ParseInt(values[defaultUpStream], 10, 64)
  402. base.DefaultUpStream = de
  403. forward, _ := strconv.ParseInt(values[forwardKey], 10, 64)
  404. var num int64
  405. for num = 256; num > 0; num /= 2 {
  406. if ((forward & num) == num) && (num != or) {
  407. base.Forward = append(base.Forward, num)
  408. }
  409. }
  410. op, _ := strconv.ParseInt(values[options], 10, 64)
  411. base.Options = op
  412. //这里判断是否有wmask mmask的流
  413. if 4&op == 4 {
  414. base.Wmask = true
  415. }
  416. if 8&op == 8 {
  417. base.Mmask = true
  418. }
  419. base.Key = values[key]
  420. if strings.Contains(n, "_bs_") {
  421. base.Type = 2
  422. } else {
  423. base.Type = 1
  424. }
  425. item.List = append(item.List, &base)
  426. }
  427. }
  428. if len(item.List) > 0 {
  429. //log.Warn("miss=%v", missRids[i])
  430. infos[missRids[i]] = &item
  431. }
  432. }
  433. }
  434. // 更新local cache
  435. d.storeMultiStreamInfo(c, infos)
  436. return infos, nil
  437. }
  438. // UpdateLastCDNCache 设置last cdn
  439. func (d *Dao) UpdateLastCDNCache(c context.Context, rid int64, origin int64) error {
  440. conn := d.redis.Get(c)
  441. defer conn.Close()
  442. key := d.getLastCDNKey(rid)
  443. if err := conn.Send("SET", key, origin); err != nil {
  444. return fmt.Errorf("redis: conn.Send(SET, %s, %v) error(%v)", key, origin, err)
  445. }
  446. if err := conn.Send("EXPIRE", key, _lastCDNExpireTime); err != nil {
  447. return fmt.Errorf("redis: conn.Send(EXPIRE key(%s) expire(%d)) error(%v)", key, _lastCDNExpireTime, err)
  448. }
  449. if err := conn.Flush(); err != nil {
  450. return fmt.Errorf("redis: conn.Flush error(%v)", err)
  451. }
  452. for i := 0; i < 2; i++ {
  453. if _, err := conn.Receive(); err != nil {
  454. return fmt.Errorf("redis: conn.Receive error(%v)", err)
  455. }
  456. }
  457. return nil
  458. }
  459. // UpdateChangeSrcCache 切流
  460. func (d *Dao) UpdateChangeSrcCache(c context.Context, rid int64, origin int64) error {
  461. conn := d.redis.Get(c)
  462. defer conn.Close()
  463. key := d.getChangeSrcKey(rid)
  464. if err := conn.Send("SET", key, origin); err != nil {
  465. return fmt.Errorf("redis: conn.Send(SET, %s, %v) error(%v)", key, origin, err)
  466. }
  467. if err := conn.Send("EXPIRE", key, _changeSrcExpireTime); err != nil {
  468. return fmt.Errorf("redis: conn.Send(EXPIRE key(%s) expire(%d)) error(%v)", key, _changeSrcExpireTime, err)
  469. }
  470. if err := conn.Flush(); err != nil {
  471. return fmt.Errorf("redis: conn.Flush error(%v)", err)
  472. }
  473. for i := 0; i < 2; i++ {
  474. if _, err := conn.Receive(); err != nil {
  475. return fmt.Errorf("redis: conn.Receive error(%v)", err)
  476. }
  477. }
  478. return nil
  479. }
  480. // UpdateStreamForwardStatus 更新forward值
  481. func (d *Dao) UpdateStreamStatusCache(c context.Context, stream *model.StreamStatus) {
  482. var (
  483. exist bool
  484. err error
  485. conn = d.redis.Get(c)
  486. allName string
  487. )
  488. defer conn.Close()
  489. // 首先判断是否存在
  490. key := d.getRoomIDKey(stream.RoomID)
  491. if exist, err = redis.Bool(conn.Do("EXISTS", key)); err != nil {
  492. if err == redis.ErrNil {
  493. return
  494. }
  495. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  496. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  497. return
  498. }
  499. if !exist {
  500. return
  501. }
  502. // 不传sname 默认为主流
  503. if stream.StreamName == "" {
  504. if allName, err = redis.String(conn.Do("HGET", key, _streamRoomFieldAllName)); err != nil {
  505. if err == redis.ErrNil {
  506. return
  507. }
  508. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  509. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  510. return
  511. }
  512. names := strings.Split(allName, "|")
  513. for _, v := range names {
  514. if !strings.Contains(v, "_bs_") {
  515. stream.StreamName = v
  516. break
  517. }
  518. }
  519. }
  520. count := 0
  521. // 是否是新增备用流
  522. if stream.Add {
  523. var names string
  524. if names, err = redis.String(conn.Do("HGET", key, _streamRoomFieldAllName)); err != nil {
  525. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  526. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  527. return
  528. }
  529. count++
  530. if err := conn.Send("HSET", key, _streamRoomFieldAllName, fmt.Sprintf("%s|%s", names, stream.StreamName)); err != nil {
  531. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  532. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  533. return
  534. }
  535. count++
  536. secret := d.getRoomFieldSecretKey(stream.StreamName)
  537. if err := conn.Send("HSET", key, secret, stream.Key); err != nil {
  538. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  539. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  540. return
  541. }
  542. }
  543. // 如果origin改变
  544. if stream.OriginChange {
  545. count++
  546. originKey := d.getRoomFieldOriginKey(stream.StreamName)
  547. if err := conn.Send("HSET", key, originKey, stream.Origin); err != nil {
  548. // 如果设置失败,则删除key
  549. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  550. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  551. return
  552. }
  553. }
  554. // forward
  555. if stream.ForwardChange {
  556. count++
  557. forwardKey := d.getRoomFieldForwardKey(stream.StreamName)
  558. if err := conn.Send("HSET", key, forwardKey, stream.Forward); err != nil {
  559. // 如果设置失败,则删除key
  560. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  561. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  562. return
  563. }
  564. }
  565. // 切上行
  566. if stream.DefaultChange {
  567. count++
  568. defaultUpKey := d.getRoomFieldDefaultKey(stream.StreamName)
  569. if err := conn.Send("HSET", key, defaultUpKey, stream.DefaultUpStream); err != nil {
  570. // 如果设置失败,则删除key
  571. log.Errorv(c, log.KV("log", fmt.Sprintf("update_status_err=%v", err)))
  572. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  573. return
  574. }
  575. }
  576. //切换options
  577. if stream.OptionsChange {
  578. if stream.Options >= 0 {
  579. count++
  580. optionsKey := d.getRoomFieldOption(stream.StreamName)
  581. if err := conn.Send("HSET", key, optionsKey, stream.Options); err != nil {
  582. // 如果设置失败,则删除key
  583. d.DeleteStreamByRIDFromCache(c, stream.RoomID)
  584. return
  585. }
  586. }
  587. }
  588. if err := conn.Flush(); err != nil {
  589. log.Infov(c, log.KV("conn.Flush error(%v)", err.Error()))
  590. return
  591. }
  592. for i := 0; i < count; i++ {
  593. if _, err := conn.Receive(); err != nil {
  594. log.Infov(c, log.KV("conn.Receive error(%v)", err.Error()))
  595. return
  596. }
  597. }
  598. }
  599. // GetLastCDNFromCache 查询上一次cdn
  600. func (d *Dao) GetLastCDNFromCache(c context.Context, rid int64) (int64, error) {
  601. conn := d.redis.Get(c)
  602. defer conn.Close()
  603. key := d.getLastCDNKey(rid)
  604. origin, err := redis.Int64(conn.Do("GET", key))
  605. if err != nil {
  606. if err != redis.ErrNil {
  607. return 0, fmt.Errorf("redis: conn.Do(GET, %s) error(%v)", key, err)
  608. }
  609. }
  610. if origin <= 0 {
  611. return 0, nil
  612. }
  613. return origin, nil
  614. }
  615. // GetChangeSrcFromCache 查询上一次cdn
  616. func (d *Dao) GetChangeSrcFromCache(c context.Context, rid int64) (int64, error) {
  617. conn := d.redis.Get(c)
  618. defer conn.Close()
  619. key := d.getChangeSrcKey(rid)
  620. origin, err := redis.Int64(conn.Do("GET", key))
  621. if err != nil {
  622. if err != redis.ErrNil {
  623. return 0, fmt.Errorf("redis: conn.Do(GET, %s) error(%v)", key, err)
  624. }
  625. }
  626. if origin <= 0 {
  627. return 0, nil
  628. }
  629. return origin, nil
  630. }
  631. // DeleteStreamByRIDFromCache 删除一个房间的的缓存信息
  632. func (d *Dao) DeleteStreamByRIDFromCache(c context.Context, rid int64) (err error) {
  633. // todo 删除内存
  634. // 删除redis
  635. // redis删除失败,进行重试,确保缓存的信息是最新的
  636. conn := d.redis.Get(c)
  637. defer conn.Close()
  638. roomKey := d.getRoomIDKey(rid)
  639. for i := 0; i < 3; i++ {
  640. _, err = conn.Do("DEL", roomKey)
  641. if err != nil {
  642. log.Error("conn.Do(DEL, %s) error(%v)", roomKey, err)
  643. continue
  644. } else {
  645. return nil
  646. }
  647. }
  648. return err
  649. }
  650. // DeleteLastCDNFromCache 删除上一次到cdn
  651. func (d *Dao) DeleteLastCDNFromCache(c context.Context, rid int64) error {
  652. conn := d.redis.Get(c)
  653. defer conn.Close()
  654. key := d.getLastCDNKey(rid)
  655. _, err := conn.Do("DEL", key)
  656. return err
  657. }
  658. // UpdateRoomOptionsCache 更新Options状态
  659. func (d *Dao) UpdateRoomOptionsCache(c context.Context, rid int64, streamname string, options int64) error {
  660. conn := d.redis.Get(c)
  661. defer conn.Close()
  662. roomKey := d.getRoomIDKey(rid)
  663. optionsKey := d.getRoomFieldOption(streamname)
  664. _, err := conn.Do("HSET", roomKey, optionsKey, options)
  665. return err
  666. }
  667. // UpdateRoomHotStatusCache 更新房间冷热流状态
  668. func (d *Dao) UpdateRoomHotStatusCache(c context.Context, rid int64, hot int64) error {
  669. conn := d.redis.Get(c)
  670. defer conn.Close()
  671. roomKey := d.getRoomIDKey(rid)
  672. hotKey := d.getRoomFieldHotKey(rid)
  673. _, err := conn.Do("HSET", roomKey, hotKey, hot)
  674. return err
  675. }