arrowTypeToDuckDB maps an Arrow DataType back to a DuckDB type name string.
(dt arrow.DataType)
| 656 | } |
| 657 | } |
| 658 | } |
| 659 | |
| 660 | // SetS3CacheMode selects the cache transport mode on the session worker. |
| 661 | func (e *FlightExecutor) SetS3CacheMode(ctx context.Context, mode string) (err error) { |
| 662 | if e.dead.Load() { |
| 663 | return ErrWorkerDead |
| 664 | } |
| 665 | if e.client == nil || e.client.Client == nil { |
| 666 | return ErrWorkerDead |
| 667 | } |
| 668 | defer recoverClientPanic(&err) |
| 669 | |
| 670 | payload, err := json.Marshal(wire.WorkerSetS3CachePayload{ |
| 671 | WorkerControlMetadata: wire.WorkerControlMetadata{ |
| 672 | WorkerID: e.workerID, |
| 673 | OwnerEpoch: e.ownerEpoch, |
| 674 | CPInstanceID: e.cpInstanceID, |
| 675 | }, |
| 676 | Mode: mode, |
| 677 | }) |
| 678 | if err != nil { |
| 679 | return err |
| 680 | } |
| 681 | |
| 682 | merged, cancel := e.mergedContext(ctx) |
| 683 | defer cancel() |
| 684 | |
| 685 | stream, err := e.client.Client.DoAction( |
| 686 | e.withSession(merged), |
| 687 | &flight.Action{Type: setSessionS3CacheModeAction, Body: payload}, |
| 688 | ) |
| 689 | if err != nil { |
| 690 | return err |
| 691 | } |
| 692 | for { |
| 693 | _, err := stream.Recv() |
| 694 | if errors.Is(err, io.EOF) { |
| 695 | return nil |
| 696 | } |
| 697 | if err != nil { |
| 698 | return err |
| 699 | } |
| 700 | } |
| 701 | } |
| 702 | |
| 703 | func (e *FlightExecutor) releaseQueryHandle(ticket *flight.Ticket) (err error) { |
| 704 | if e.dead.Load() { |
| 705 | return nil |
| 706 | } |
| 707 | if e.client == nil || e.client.Client == nil || ticket == nil || len(ticket.Ticket) == 0 { |
| 708 | return nil |
| 709 | } |
| 710 | defer func() { |
| 711 | recoverClientPanic(&err) |
| 712 | if err != nil && isTerminalSessionIdleWaitError(err) { |
| 713 | err = nil |
| 714 | } |
| 715 | }() |