(ctx context.Context, waitCount int, maxWaitTime time.Duration)
| 361 | } |
| 362 | |
| 363 | func (em *EventManager) waitForProcessedTotal(ctx context.Context, waitCount int, maxWaitTime time.Duration) error { |
| 364 | startTime := time.Now() |
| 365 | |
| 366 | worker := func() (bool, error, interface{}) { |
| 367 | eventTotal := em.GetEventsProcessedSuccess() + em.GetEventsProcessedFail() |
| 368 | if eventTotal >= int64(waitCount) { |
| 369 | base.DebugfCtx(ctx, base.KeyAll, "waitForProcessedTotal(%d) took %v", waitCount, time.Since(startTime)) |
| 370 | return false, nil, nil |
| 371 | } |
| 372 | |
| 373 | return true, nil, nil |
| 374 | } |
| 375 | |
| 376 | ctx, cancel := context.WithDeadline(ctx, startTime.Add(maxWaitTime)) |
| 377 | sleeper := base.SleeperFuncCtx(base.CreateMaxDoublingSleeperFunc(math.MaxInt64, 1, 1000), ctx) |
| 378 | err, _ := base.RetryLoop(ctx, fmt.Sprintf("waitForProcessedTotal(%d)", waitCount), worker, sleeper) |
| 379 | cancel() |
| 380 | return err |
| 381 | } |
| 382 | |
| 383 | func GetRouterWithHandler(wr *WebhookRequest) http.Handler { |
| 384 | r := http.NewServeMux() |
no test coverage detected