mirror of
https://github.com/ScuroNeko/mtg.git
synced 2026-08-31 09:54:01 +03:00
Rewrite to WaitGroup.Go
This commit is contained in:
+27
-54
@@ -12,14 +12,11 @@ type multiObserver struct {
|
|||||||
|
|
||||||
func (m multiObserver) EventStart(evt mtglib.EventStart) {
|
func (m multiObserver) EventStart(evt mtglib.EventStart) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventStart(evt)
|
||||||
|
})
|
||||||
obs.EventStart(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -27,14 +24,11 @@ func (m multiObserver) EventStart(evt mtglib.EventStart) {
|
|||||||
|
|
||||||
func (m multiObserver) EventConnectedToDC(evt mtglib.EventConnectedToDC) {
|
func (m multiObserver) EventConnectedToDC(evt mtglib.EventConnectedToDC) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventConnectedToDC(evt)
|
||||||
|
})
|
||||||
obs.EventConnectedToDC(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -42,14 +36,11 @@ func (m multiObserver) EventConnectedToDC(evt mtglib.EventConnectedToDC) {
|
|||||||
|
|
||||||
func (m multiObserver) EventDomainFronting(evt mtglib.EventDomainFronting) {
|
func (m multiObserver) EventDomainFronting(evt mtglib.EventDomainFronting) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventDomainFronting(evt)
|
||||||
|
})
|
||||||
obs.EventDomainFronting(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -57,14 +48,11 @@ func (m multiObserver) EventDomainFronting(evt mtglib.EventDomainFronting) {
|
|||||||
|
|
||||||
func (m multiObserver) EventTraffic(evt mtglib.EventTraffic) {
|
func (m multiObserver) EventTraffic(evt mtglib.EventTraffic) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventTraffic(evt)
|
||||||
|
})
|
||||||
obs.EventTraffic(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -72,14 +60,11 @@ func (m multiObserver) EventTraffic(evt mtglib.EventTraffic) {
|
|||||||
|
|
||||||
func (m multiObserver) EventFinish(evt mtglib.EventFinish) {
|
func (m multiObserver) EventFinish(evt mtglib.EventFinish) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventFinish(evt)
|
||||||
|
})
|
||||||
obs.EventFinish(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -87,14 +72,11 @@ func (m multiObserver) EventFinish(evt mtglib.EventFinish) {
|
|||||||
|
|
||||||
func (m multiObserver) EventConcurrencyLimited(evt mtglib.EventConcurrencyLimited) {
|
func (m multiObserver) EventConcurrencyLimited(evt mtglib.EventConcurrencyLimited) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventConcurrencyLimited(evt)
|
||||||
|
})
|
||||||
obs.EventConcurrencyLimited(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -102,14 +84,11 @@ func (m multiObserver) EventConcurrencyLimited(evt mtglib.EventConcurrencyLimite
|
|||||||
|
|
||||||
func (m multiObserver) EventIPBlocklisted(evt mtglib.EventIPBlocklisted) {
|
func (m multiObserver) EventIPBlocklisted(evt mtglib.EventIPBlocklisted) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventIPBlocklisted(evt)
|
||||||
|
})
|
||||||
obs.EventIPBlocklisted(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -117,14 +96,11 @@ func (m multiObserver) EventIPBlocklisted(evt mtglib.EventIPBlocklisted) {
|
|||||||
|
|
||||||
func (m multiObserver) EventReplayAttack(evt mtglib.EventReplayAttack) {
|
func (m multiObserver) EventReplayAttack(evt mtglib.EventReplayAttack) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventReplayAttack(evt)
|
||||||
|
})
|
||||||
obs.EventReplayAttack(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
@@ -132,14 +108,11 @@ func (m multiObserver) EventReplayAttack(evt mtglib.EventReplayAttack) {
|
|||||||
|
|
||||||
func (m multiObserver) EventIPListSize(evt mtglib.EventIPListSize) {
|
func (m multiObserver) EventIPListSize(evt mtglib.EventIPListSize) {
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(m.observers))
|
|
||||||
|
|
||||||
for _, v := range m.observers {
|
for _, v := range m.observers {
|
||||||
go func(obs Observer) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
v.EventIPListSize(evt)
|
||||||
|
})
|
||||||
obs.EventIPListSize(evt)
|
|
||||||
}(v)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|||||||
+4
-10
@@ -61,11 +61,8 @@ func (a *Access) Run(cli *CLI, version string) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(2)
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
|
wg.Go(func() {
|
||||||
ip := a.PublicIPv4
|
ip := a.PublicIPv4
|
||||||
if ip == nil {
|
if ip == nil {
|
||||||
ip = a.getIP(ntw, "tcp4")
|
ip = a.getIP(ntw, "tcp4")
|
||||||
@@ -76,11 +73,8 @@ func (a *Access) Run(cli *CLI, version string) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
resp.IPv4 = a.makeURLs(conf, ip)
|
resp.IPv4 = a.makeURLs(conf, ip)
|
||||||
}()
|
})
|
||||||
|
wg.Go(func() {
|
||||||
go func() {
|
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
ip := a.PublicIPv6
|
ip := a.PublicIPv6
|
||||||
if ip == nil {
|
if ip == nil {
|
||||||
ip = a.getIP(ntw, "tcp6")
|
ip = a.getIP(ntw, "tcp6")
|
||||||
@@ -91,7 +85,7 @@ func (a *Access) Run(cli *CLI, version string) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
resp.IPv6 = a.makeURLs(conf, ip)
|
resp.IPv6 = a.makeURLs(conf, ip)
|
||||||
}()
|
})
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
|
|||||||
@@ -112,18 +112,15 @@ func (f *Firehol) update() {
|
|||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(len(f.blocklists))
|
|
||||||
|
|
||||||
mutex := &sync.Mutex{}
|
mutex := &sync.Mutex{}
|
||||||
ranger := cidranger.NewPCTrieRanger()
|
ranger := cidranger.NewPCTrieRanger()
|
||||||
|
|
||||||
for _, v := range f.blocklists {
|
for _, v := range f.blocklists {
|
||||||
go func(file files.File) {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
logger := f.logger.BindStr("filename", v.String())
|
||||||
|
|
||||||
logger := f.logger.BindStr("filename", file.String())
|
fileContent, err := v.Open(ctx)
|
||||||
|
|
||||||
fileContent, err := file.Open(ctx)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.WarningError("update has failed", err)
|
logger.WarningError("update has failed", err)
|
||||||
|
|
||||||
@@ -135,7 +132,7 @@ func (f *Firehol) update() {
|
|||||||
if err := f.updateFromFile(mutex, ranger, bufio.NewScanner(fileContent)); err != nil {
|
if err := f.updateFromFile(mutex, ranger, bufio.NewScanner(fileContent)); err != nil {
|
||||||
logger.WarningError("update has failed", err)
|
logger.WarningError("update has failed", err)
|
||||||
}
|
}
|
||||||
}(v)
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|||||||
@@ -52,17 +52,9 @@ func (suite *CircuitBreakerTestSuite) TestMultipleRunsOk() {
|
|||||||
Return(suite.connMock, nil)
|
Return(suite.connMock, nil)
|
||||||
|
|
||||||
wg := &sync.WaitGroup{}
|
wg := &sync.WaitGroup{}
|
||||||
wg.Add(5)
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
wg.Wait()
|
|
||||||
suite.ctxCancel()
|
|
||||||
}()
|
|
||||||
|
|
||||||
for range 5 {
|
for range 5 {
|
||||||
go func() {
|
wg.Go(func() {
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
conn, err := suite.d.DialContext(suite.ctx, "tcp", "127.0.0.1")
|
conn, err := suite.d.DialContext(suite.ctx, "tcp", "127.0.0.1")
|
||||||
|
|
||||||
suite.mutex.Lock()
|
suite.mutex.Lock()
|
||||||
@@ -70,9 +62,14 @@ func (suite *CircuitBreakerTestSuite) TestMultipleRunsOk() {
|
|||||||
|
|
||||||
suite.NoError(err)
|
suite.NoError(err)
|
||||||
suite.Equal("127.0.0.1:3128", conn.RemoteAddr().String())
|
suite.Equal("127.0.0.1:3128", conn.RemoteAddr().String())
|
||||||
}()
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
wg.Wait()
|
||||||
|
suite.ctxCancel()
|
||||||
|
}()
|
||||||
|
|
||||||
suite.Eventually(func() bool {
|
suite.Eventually(func() bool {
|
||||||
_, ok := <-suite.ctx.Done()
|
_, ok := <-suite.ctx.Done()
|
||||||
|
|
||||||
|
|||||||
+4
-12
@@ -81,32 +81,24 @@ func (n *network) dnsResolve(protocol, address string) ([]string, error) {
|
|||||||
|
|
||||||
switch protocol {
|
switch protocol {
|
||||||
case "tcp", "tcp4":
|
case "tcp", "tcp4":
|
||||||
wg.Add(1)
|
wg.Go(func() {
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
resolved := n.dns.LookupA(address)
|
resolved := n.dns.LookupA(address)
|
||||||
|
|
||||||
mutex.Lock()
|
mutex.Lock()
|
||||||
ips = append(ips, resolved...)
|
ips = append(ips, resolved...)
|
||||||
mutex.Unlock()
|
mutex.Unlock()
|
||||||
}()
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
switch protocol {
|
switch protocol {
|
||||||
case "tcp", "tcp6":
|
case "tcp", "tcp6":
|
||||||
wg.Add(1)
|
wg.Go(func() {
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
resolved := n.dns.LookupAAAA(address)
|
resolved := n.dns.LookupAAAA(address)
|
||||||
|
|
||||||
mutex.Lock()
|
mutex.Lock()
|
||||||
ips = append(ips, resolved...)
|
ips = append(ips, resolved...)
|
||||||
mutex.Unlock()
|
mutex.Unlock()
|
||||||
}()
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|||||||
Reference in New Issue
Block a user