Revert "Add broadcast stats"

This reverts commit 8b096dbeb8.
This commit is contained in:
Sammy Libre 2017-02-02 19:50:06 +05:00
parent acb710f532
commit 8ea446f333

View File

@ -4,9 +4,7 @@ import (
"log" "log"
"regexp" "regexp"
"strings" "strings"
"sync"
"sync/atomic" "sync/atomic"
"time"
"../util" "../util"
) )
@ -108,34 +106,24 @@ func (s *StratumServer) broadcastNewJobs() {
defer s.sessionsMu.RUnlock() defer s.sessionsMu.RUnlock()
count := len(s.sessions) count := len(s.sessions)
log.Printf("Broadcasting new jobs to %d miners", count) log.Printf("Broadcasting new jobs to %d miners", count)
start := time.Now() bcast := make(chan int, 1024*16)
slots := make(chan bool, 1024*16) n := 0
var ok, fails int64
var wg sync.WaitGroup
for m := range s.sessions { for m := range s.sessions {
wg.Add(1) n++
slots <- true bcast <- n
go func(cs *Session) { go func(cs *Session) {
reply := cs.getJob(t) reply := cs.getJob(t)
err := cs.pushMessage("job", &reply) err := cs.pushMessage("job", &reply)
<-slots <-bcast
if err != nil { if err != nil {
log.Printf("Job transmit error to %s: %v", cs.ip, err) log.Printf("Job transmit error to %s: %v", cs.ip, err)
atomic.AddInt64(&fails, 1)
wg.Done()
s.removeSession(cs) s.removeSession(cs)
} else { } else {
atomic.AddInt64(&ok, 1)
s.setDeadline(cs.conn) s.setDeadline(cs.conn)
wg.Done()
} }
}(m) }(m)
} }
wg.Wait()
log.Printf("Done jobs broadcast in %s for %d/%d/%d miners", time.Since(start), count, ok, fails)
} }
func (s *StratumServer) refreshBlockTemplate(bcast bool) { func (s *StratumServer) refreshBlockTemplate(bcast bool) {