waitForWorker polls for the worker socket and creates a Flight SQL client.
(socketPath, bearerToken string, timeout time.Duration)
| 514 | p.mu.Lock() |
| 515 | p.workers[id] = w |
| 516 | workerCount := len(p.workers) |
| 517 | p.mu.Unlock() |
| 518 | observeControlPlaneWorkers(workerCount) |
| 519 | |
| 520 | slog.Info("Worker spawned.", "id", id, "pid", cmd.Process.Pid, "socket", socketPath) |
| 521 | return nil |
| 522 | } |
| 523 | |
| 524 | // waitForWorker polls for the worker socket and creates a Flight SQL client. |
| 525 | func waitForWorker(socketPath, bearerToken string, timeout time.Duration) (*flightsql.Client, error) { |
| 526 | deadline := time.Now().Add(timeout) |
| 527 | var lastErr error |
| 528 | attempts := 0 |
| 529 | |
| 530 | for time.Now().Before(deadline) { |
| 531 | if _, err := os.Stat(socketPath); err == nil { |
| 532 | // Socket exists, try to connect |
| 533 | addr := "unix://" + socketPath |
| 534 | var dialOpts []grpc.DialOption |
| 535 | dialOpts = append(dialOpts, grpc.WithTransportCredentials(insecure.NewCredentials())) |
| 536 | dialOpts = append(dialOpts, grpc.WithDefaultCallOptions( |
| 537 | grpc.MaxCallRecvMsgSize(flightclient.MaxGRPCMessageSize), |
| 538 | grpc.MaxCallSendMsgSize(flightclient.MaxGRPCMessageSize), |
| 539 | )) |
| 540 | |
| 541 | if bearerToken != "" { |
| 542 | dialOpts = append(dialOpts, grpc.WithPerRPCCredentials(&workerBearerCreds{token: bearerToken})) |
| 543 | } |
| 544 | |
| 545 | client, err := flightsql.NewClient(addr, nil, nil, dialOpts...) |
| 546 | if err == nil { |
| 547 | // Verify with a health check |
| 548 | ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) |
| 549 | _, err = doHealthCheck(ctx, client) |
| 550 | cancel() |
| 551 | if err == nil { |
| 552 | return client, nil |
| 553 | } |
| 554 | lastErr = err |
| 555 | _ = client.Close() |
| 556 | } else { |
| 557 | lastErr = fmt.Errorf("grpc dial: %w", err) |
| 558 | } |
| 559 | attempts++ |
| 560 | if attempts <= 3 || attempts%10 == 0 { |
| 561 | slog.Debug("waitForWorker health check attempt failed.", "socket", socketPath, "attempt", attempts, "error", lastErr) |
| 562 | } |
| 563 | } else { |
| 564 | lastErr = err |
no test coverage detected