| 883 | return |
| 884 | } |
| 885 | } |
| 886 | }() |
| 887 | |
| 888 | return schema, ch, nil |
| 889 | } |
| 890 | |
| 891 | func (h *FlightSQLHandler) DoPutCommandStatementUpdate(ctx context.Context, |
| 892 | cmd flightsql.StatementUpdate) (int64, error) { |
| 893 | |
| 894 | session, err := h.sessionFromContext(ctx) |
| 895 | if err != nil { |
| 896 | return 0, err |
| 897 | } |
| 898 | finishOperation, ok := session.beginOperation() |
| 899 | if busyErr := sessionBusyStatus(ok); busyErr != nil { |
| 900 | return 0, busyErr |
| 901 | } |
| 902 | session.setCurrentQueryID(queryIDFromContext(ctx)) |
| 903 | defer func() { |
| 904 | session.setCurrentQueryID("") |
| 905 | finishOperation() |
| 906 | }() |
| 907 | |
| 908 | tx, _, ttx, err := session.getOpenTxn(cmd.GetTransactionId()) |
| 909 | if err != nil { |
| 910 | return 0, err |
| 911 | } |
| 912 | |
| 913 | query := cmd.GetQuery() |
| 914 | var endConnWork func() |
| 915 | if tx == nil && !isEmptyFlightQuery(query) { |
| 916 | var ok bool |
| 917 | endConnWork, ok = session.beginConnWork() |
| 918 | if !ok { |
| 919 | return 0, status.Error(codes.NotFound, "session closed") |
| 920 | } |
| 921 | defer endConnWork() |
| 922 | } |
| 923 | finishDrain, err := h.pool.beginDrainWork(session.allowsDrainContinuation(string(cmd.GetTransactionId()))) |
| 924 | if drainErr := workerDrainingStatus(err); drainErr != nil { |
| 925 | return 0, drainErr |
| 926 | } |
| 927 | if err != nil { |
| 928 | return 0, status.Errorf(codes.Internal, "start update drain tracking: %v", err) |
| 929 | } |