(router IRoutes, pubSub *p2p.PubSub)
| 49 | } |
| 50 | |
| 51 | func subscribe(router IRoutes, pubSub *p2p.PubSub) { |
| 52 | utils.GoWithRecover(func() { |
| 53 | for args := range pubSub.Inbound { |
| 54 | // logrus.Info("inbound: ", args.Message) |
| 55 | var data []interface{} |
| 56 | err := json.Unmarshal([]byte(args.Message), &data) |
| 57 | if err != nil { |
| 58 | logrus.Errorf("subscribe error: %v", err) |
| 59 | continue |
| 60 | } |
| 61 | err = router.Sync(data) |
| 62 | if err != nil { |
| 63 | logrus.Errorf("subscribe sync error: %v", err) |
| 64 | } |
| 65 | } |
| 66 | }, func(r interface{}) { |
| 67 | time.Sleep(time.Second) |
| 68 | subscribe(router, pubSub) |
| 69 | }) |
| 70 | } |
no test coverage detected