func (c *Client) Send(ctx context.Context, messages []Message, maxWorkers int) error {
sem := make(chan struct{}, maxWorkers)
var wg sync.WaitGroup
var resultErr error
var once sync.Once
for _, msg := range messages {
wg.Add(1)
go func(msg Message) {
defer wg.Done()
sem <- struct{}{}
defer func() {
<-sem
}()
body, err := json.Marshal(msg)
if err != nil {
once.Do(func() { resultErr = err; cancel() })
return
}
req, _ := http.NewRequestWithContext(
ctx,
http.MethodPost,
c.url,
bytes.NewReader(body),
)
req.Header.Set("Content-Type", "application/json")
resp, err := c.httpClient.Do(req)
if err != nil {
once.Do(func() { resultErr = err; cancel()})
return
}
defer resp.Body.Close()
if resp.StatusCode >= 500 {
once.Do(func() {
resultErr = fmt.Errorf("server error")
cancel()})
return
}
c.sentCount.Add(1)
time.Sleep(time.Second)
}(msg)
}
wg.Wait()
}