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

Method BeginTransaction

duckdbservice/flight_handler.go:885–926  ·  view source on GitHub ↗
(ctx context.Context,
	req flightsql.ActionBeginTransactionRequest)

Source from the content-addressed store, hash-verified

883 return
884 }
885 }
886 }()
887
888 return schema, ch, nil
889}
890
891func (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 }

Calls 12

sessionFromContextMethod · 0.95
sessionBusyStatusFunction · 0.85
workerDrainingStatusFunction · 0.85
sessionClosedStatusFunction · 0.85
addTrackedTransactionFunction · 0.85
beginOperationMethod · 0.80
beginDrainWorkMethod · 0.80
beginTxMethod · 0.80
AddMethod · 0.80
NowMethod · 0.80
rollbackTxMethod · 0.80
ErrorMethod · 0.45