(frontendContext context.Context, frontendAddr string, msg *schedulerpb.FrontendToScheduler, fragment plan_fragments.Fragment)
| 383 | } |
| 384 | |
| 385 | func (s *Scheduler) enqueueRequest(frontendContext context.Context, frontendAddr string, msg *schedulerpb.FrontendToScheduler, fragment plan_fragments.Fragment) error { |
| 386 | // Create new context for this request, to support cancellation. |
| 387 | ctx, cancel := context.WithCancel(frontendContext) |
| 388 | shouldCancel := true |
| 389 | defer func() { |
| 390 | if shouldCancel { |
| 391 | cancel() |
| 392 | } |
| 393 | }() |
| 394 | |
| 395 | // Extract tracing information from headers in HTTP request. FrontendContext doesn't have the correct tracing |
| 396 | // information, since that is a long-running request. |
| 397 | tracer := opentracing.GlobalTracer() |
| 398 | parentSpanContext, err := httpgrpcutil.GetParentSpanForRequest(tracer, msg.HttpRequest) |
| 399 | if err != nil { |
| 400 | return err |
| 401 | } |
| 402 | |
| 403 | userID := msg.GetUserID() |
| 404 | |
| 405 | req := &schedulerRequest{ |
| 406 | frontendAddress: frontendAddr, |
| 407 | userID: msg.UserID, |
| 408 | queryID: msg.QueryID, |
| 409 | request: msg.HttpRequest, |
| 410 | statsEnabled: msg.StatsEnabled, |
| 411 | fragment: fragment, |
| 412 | } |
| 413 | |
| 414 | now := time.Now() |
| 415 | |
| 416 | req.parentSpanContext = parentSpanContext |
| 417 | req.queueSpan, req.ctx = opentracing.StartSpanFromContextWithTracer(ctx, tracer, "queued", opentracing.ChildOf(parentSpanContext)) |
| 418 | req.enqueueTime = now |
| 419 | req.ctxCancel = cancel |
| 420 | |
| 421 | // aggregate the max queriers limit in the case of a multi tenant query |
| 422 | tenantIDs, err := users.TenantIDsFromOrgID(userID) |
| 423 | if err != nil { |
| 424 | return err |
| 425 | } |
| 426 | maxQueriers := validation.SmallestPositiveNonZeroFloat64PerTenant(tenantIDs, s.limits.MaxQueriersPerUser) |
| 427 | |
| 428 | s.activeUsers.UpdateUserTimestamp(userID, now) |
| 429 | return s.requestQueue.EnqueueRequest(userID, req, maxQueriers, func() { |
| 430 | shouldCancel = false |
| 431 | |
| 432 | s.pendingRequestsMu.Lock() |
| 433 | defer s.pendingRequestsMu.Unlock() |
| 434 | |
| 435 | queryKey := queryKey{frontendAddr: frontendAddr, queryID: msg.QueryID} |
| 436 | s.queryFragmentRegistry[queryKey] = append(s.queryFragmentRegistry[queryKey], req.fragment.FragmentID) |
| 437 | s.pendingRequests[requestKey{queryKey: queryKey, fragmentID: req.fragment.FragmentID}] = req |
| 438 | }) |
| 439 | } |
| 440 | |
| 441 | // This method doesn't do removal from the queue. |
| 442 | func (s *Scheduler) cancelRequestAndRemoveFromPending(frontendAddr string, queryID uint64, fragmentID uint64, cancelAll bool) { |
no test coverage detected