MCPcopy Create free account
hub / github.com/PostHog/duckgres / waitForSessionIdle

Method waitForSessionIdle

server/flightclient/flight_executor.go:349–408  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

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
361func (e *FlightExecutor) Query(query string, args ...any) (sqlcore.RowSet, error) {
362 return e.QueryContext(context.Background(), query, args...)
363}
364
365func (e *FlightExecutor) Exec(query string, args ...any) (sqlcore.ExecResult, error) {
366 return e.ExecContext(context.Background(), query, args...)
367}
368
369func (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
373func (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).
386func (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
400func (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 {

Callers 1

QueryContextMethod · 0.95

Calls 5

withSessionMethod · 0.95
recoverClientPanicFunction · 0.85
DoActionMethod · 0.45
RecvMethod · 0.45

Tested by

no test coverage detected