()
| 347 | // ExecuteUpdate is a bidi-streaming RPC and grpc.Trailer(&md) as a |
| 348 | // CallOption only works for unary RPCs — for streams the trailer is |
| 349 | // only reachable via stream.Trailer() after Recv returns io.EOF. |
| 350 | // Without that capture the worker's per-query profiling JSON is |
| 351 | // silently dropped on this side. See executeUpdateWithTrailer. |
| 352 | affected, trailer, err := executeUpdateWithTrailer(reqCtx, e.client, query) |
| 353 | e.storeProfilingFromTrailer(trailer) |
| 354 | if err != nil { |
| 355 | return nil, fmt.Errorf("flight execute update: %w", err) |
| 356 | } |
| 357 | |
| 358 | return &flightExecResult{rowsAffected: affected}, nil |
| 359 | } |
| 360 | |
| 361 | func (e *FlightExecutor) Query(query string, args ...any) (sqlcore.RowSet, error) { |
| 362 | return e.QueryContext(context.Background(), query, args...) |
| 363 | } |
| 364 | |
| 365 | func (e *FlightExecutor) Exec(query string, args ...any) (sqlcore.ExecResult, error) { |
| 366 | return e.ExecContext(context.Background(), query, args...) |
| 367 | } |
| 368 | |
| 369 | func (e *FlightExecutor) ConnContext(ctx context.Context) (sqlcore.RawConn, error) { |
| 370 | return nil, fmt.Errorf("ConnContext not supported in Flight mode (use batched INSERT for COPY FROM)") |
| 371 | } |
| 372 | |
| 373 | func (e *FlightExecutor) PingContext(ctx context.Context) error { |
| 374 | // Use a simple query to verify connectivity |
| 375 | rows, err := e.QueryContext(ctx, "SELECT 1") |
| 376 | if err != nil { |
| 377 | return fmt.Errorf("flight ping: %w", err) |
| 378 | } |
| 379 | return rows.Close() |
| 380 | } |
| 381 | |
| 382 | // mergedContext returns a context that is cancelled when either the caller's |
| 383 | // context or the executor's base context is done. This ensures gRPC calls are |
| 384 | // cancelled both when the client disconnects (caller ctx) and when the |
| 385 | // executor is closed (e.g. worker crash). |
| 386 | func (e *FlightExecutor) mergedContext(ctx context.Context) (context.Context, context.CancelFunc) { |
| 387 | merged, cancel := context.WithCancel(ctx) |
| 388 | if e.ctx != nil { |
| 389 | go func() { |
| 390 | select { |
| 391 | case <-e.ctx.Done(): |
| 392 | cancel() |
| 393 | case <-merged.Done(): |
| 394 | } |
| 395 | }() |
| 396 | } |
| 397 | return merged, cancel |
| 398 | } |
| 399 | |
| 400 | func (e *FlightExecutor) Close() error { |
| 401 | e.stopQueryLogForwarding() |
| 402 | if !e.waitQueryLogForwarding() && e.cancel != nil { |
| 403 | e.cancel() |
| 404 | e.queryLogWG.Wait() |
| 405 | } |
| 406 | if e.cancel != nil { |
no test coverage detected