Flight SQL method implementations
(ctx context.Context, cmd flightsql.StatementQuery, desc *flight.FlightDescriptor)
| 496 | return status.Errorf(codes.Canceled, "wait for session idle: %v", err) |
| 497 | case errors.Is(err, context.DeadlineExceeded): |
| 498 | return status.Errorf(codes.DeadlineExceeded, "wait for session idle: %v", err) |
| 499 | default: |
| 500 | return status.Errorf(codes.Internal, "wait for session idle: %v", err) |
| 501 | } |
| 502 | } |
| 503 | |
| 504 | resp, _ := json.Marshal(map[string]bool{"ok": true}) |
| 505 | return sendActionResult(stream, &flight.Result{Body: resp}) |
| 506 | } |
| 507 | |
| 508 | // doSetSessionS3Cache applies the `duckgres.s3_cache` session GUC: it swaps |
| 509 | // the tenant S3 secret between the cache-proxy transport (enabled) and the |
| 510 | // org's native HTTPS transport (bypassed). Requires a live session — the swap |
| 511 | // is instance-global, which is session-scoped only under the one-session-per- |
| 512 | // worker contract, and a sessionless caller has no business toggling it. |
| 513 | // Errors surface to the control plane so the client's SET fails rather than |
| 514 | // silently not taking effect. |
| 515 | func (h *FlightSQLHandler) doSetSessionS3Cache(body []byte, stream flight.FlightService_DoActionServer) error { |
| 516 | var req server.WorkerSetS3CachePayload |
| 517 | if err := json.Unmarshal(body, &req); err != nil { |
| 518 | return status.Errorf(codes.InvalidArgument, "invalid SetSessionS3Cache request: %v", err) |
| 519 | } |
| 520 | if err := h.pool.validateControlMetadata(req.WorkerControlMetadata); err != nil { |
| 521 | return status.Errorf(codes.FailedPrecondition, "stale worker owner: %v", err) |
| 522 | } |
| 523 | if _, err := h.sessionFromContext(stream.Context()); err != nil { |
| 524 | return err |
| 525 | } |
| 526 | mode := req.Mode |
| 527 | if mode == "" { |
| 528 | mode = "off" |
| 529 | if req.Enabled { |
| 530 | mode = "on" |
| 531 | } |
| 532 | } |
| 533 | if err := h.pool.SetS3CacheMode(mode); err != nil { |
| 534 | return status.Errorf(codes.Internal, "set session s3 cache: %v", err) |
| 535 | } |
| 536 | |
| 537 | resp, _ := json.Marshal(map[string]bool{"ok": true}) |
| 538 | return sendActionResult(stream, &flight.Result{Body: resp}) |
| 539 | } |
| 540 | |
| 541 | func (h *FlightSQLHandler) doReleaseQueryHandle(body []byte, stream flight.FlightService_DoActionServer) error { |
| 542 | var req server.WorkerReleaseQueryHandlePayload |
| 543 | if err := json.Unmarshal(body, &req); err != nil { |
| 544 | return status.Errorf(codes.InvalidArgument, "invalid ReleaseQueryHandle request: %v", err) |
| 545 | } |
| 546 | if err := h.pool.validateControlMetadata(req.WorkerControlMetadata); err != nil { |
| 547 | return status.Errorf(codes.FailedPrecondition, "stale worker owner: %v", err) |
| 548 | } |
| 549 | session, err := h.sessionFromContext(stream.Context()) |
| 550 | if err != nil { |
| 551 | return err |
| 552 | } |
| 553 | ticket, err := flightsql.GetStatementQueryTicket(&flight.Ticket{Ticket: req.Ticket}) |
| 554 | if err != nil { |
| 555 | return status.Errorf(codes.InvalidArgument, "invalid statement ticket: %v", err) |