You can not select more than 25 topics Topics must start with a chinese character,a letter or number, can include dashes ('-') and can be up to 35 characters long.

executor_at.go 5.1 kB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200
  1. /*
  2. * Licensed to the Apache Software Foundation (ASF) under one or more
  3. * contributor license agreements. See the NOTICE file distributed with
  4. * this work for additional information regarding copyright ownership.
  5. * The ASF licenses this file to You under the Apache License, Version 2.0
  6. * (the "License"); you may not use this file except in compliance with
  7. * the License. You may obtain a copy of the License at
  8. *
  9. * http://www.apache.org/licenses/LICENSE-2.0
  10. *
  11. * Unless required by applicable law or agreed to in writing, software
  12. * distributed under the License is distributed on an "AS IS" BASIS,
  13. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  14. * See the License for the specific language governing permissions and
  15. * limitations under the License.
  16. */
  17. package at
  18. import (
  19. "context"
  20. "fmt"
  21. "github.com/mitchellh/copystructure"
  22. "github.com/pkg/errors"
  23. "github.com/seata/seata-go/pkg/datasource/sql/parser"
  24. "github.com/seata/seata-go/pkg/datasource/sql/undo"
  25. "github.com/seata/seata-go/pkg/tm"
  26. "github.com/seata/seata-go/pkg/datasource/sql/exec"
  27. "github.com/seata/seata-go/pkg/datasource/sql/types"
  28. )
  29. type ATExecutor struct {
  30. hooks []exec.SQLHook
  31. ex exec.SQLExecutor
  32. }
  33. func (e *ATExecutor) Interceptors(hooks []exec.SQLHook) {
  34. e.hooks = hooks
  35. }
  36. func (e *ATExecutor) ExecWithNamedValue(ctx context.Context, execCtx *types.ExecContext, f exec.CallbackWithNamedValue) (types.ExecResult, error) {
  37. for _, hook := range e.hooks {
  38. hook.Before(ctx, execCtx)
  39. }
  40. var (
  41. beforeImages []*types.RecordImage
  42. afterImages []*types.RecordImage
  43. result types.ExecResult
  44. err error
  45. )
  46. beforeImages, err = e.beforeImage(ctx, execCtx)
  47. if err != nil {
  48. return nil, err
  49. }
  50. if beforeImages != nil {
  51. beforeImagesTmp, err := copystructure.Copy(beforeImages)
  52. if err != nil {
  53. return nil, err
  54. }
  55. newBeforeImages, ok := beforeImagesTmp.([]*types.RecordImage)
  56. if !ok {
  57. return nil, errors.New("copy beforeImages failed")
  58. }
  59. execCtx.TxCtx.RoundImages.AppendBeofreImages(newBeforeImages)
  60. }
  61. defer func() {
  62. for _, hook := range e.hooks {
  63. hook.After(ctx, execCtx)
  64. }
  65. }()
  66. if e.ex != nil {
  67. result, err = e.ex.ExecWithNamedValue(ctx, execCtx, f)
  68. } else {
  69. result, err = f(ctx, execCtx.Query, execCtx.NamedValues)
  70. }
  71. if err != nil {
  72. return nil, err
  73. }
  74. afterImages, err = e.afterImage(ctx, execCtx, beforeImages)
  75. if err != nil {
  76. return nil, err
  77. }
  78. if afterImages != nil {
  79. execCtx.TxCtx.RoundImages.AppendAfterImages(afterImages)
  80. }
  81. return result, err
  82. }
  83. func (e *ATExecutor) prepareUndoLog(ctx context.Context, execCtx *types.ExecContext) error {
  84. if execCtx.TxCtx.RoundImages.IsEmpty() {
  85. return nil
  86. }
  87. if execCtx.ParseContext.UpdateStmt != nil {
  88. if !execCtx.TxCtx.RoundImages.IsBeforeAfterSizeEq() {
  89. return fmt.Errorf("Before image size is not equaled to after image size, probably because you updated the primary keys.")
  90. }
  91. }
  92. undoLogManager, err := undo.GetUndoLogManager(execCtx.DBType)
  93. if err != nil {
  94. return err
  95. }
  96. return undoLogManager.FlushUndoLog(execCtx.TxCtx, execCtx.Conn)
  97. }
  98. func (e *ATExecutor) ExecWithValue(ctx context.Context, execCtx *types.ExecContext, f exec.CallbackWithValue) (types.ExecResult, error) {
  99. for _, hook := range e.hooks {
  100. hook.Before(ctx, execCtx)
  101. }
  102. var (
  103. beforeImages []*types.RecordImage
  104. afterImages []*types.RecordImage
  105. result types.ExecResult
  106. err error
  107. )
  108. beforeImages, err = e.beforeImage(ctx, execCtx)
  109. if err != nil {
  110. return nil, err
  111. }
  112. if beforeImages != nil {
  113. execCtx.TxCtx.RoundImages.AppendBeofreImages(beforeImages)
  114. }
  115. defer func() {
  116. for _, hook := range e.hooks {
  117. hook.After(ctx, execCtx)
  118. }
  119. }()
  120. if e.ex != nil {
  121. result, err = e.ex.ExecWithValue(ctx, execCtx, f)
  122. } else {
  123. result, err = f(ctx, execCtx.Query, execCtx.Values)
  124. }
  125. if err != nil {
  126. return nil, err
  127. }
  128. afterImages, err = e.afterImage(ctx, execCtx, beforeImages)
  129. if err != nil {
  130. return nil, err
  131. }
  132. if afterImages != nil {
  133. execCtx.TxCtx.RoundImages.AppendAfterImages(afterImages)
  134. }
  135. return result, err
  136. }
  137. func (e *ATExecutor) beforeImage(ctx context.Context, execCtx *types.ExecContext) ([]*types.RecordImage, error) {
  138. if !tm.IsGlobalTx(ctx) {
  139. return nil, nil
  140. }
  141. pc, err := parser.DoParser(execCtx.Query)
  142. if err != nil {
  143. return nil, err
  144. }
  145. if !pc.HasValidStmt() {
  146. return nil, nil
  147. }
  148. execCtx.ParseContext = pc
  149. builder := undo.GetUndologBuilder(pc.ExecutorType)
  150. if builder == nil {
  151. return nil, nil
  152. }
  153. return builder.BeforeImage(ctx, execCtx)
  154. }
  155. // After
  156. func (e *ATExecutor) afterImage(ctx context.Context, execCtx *types.ExecContext, beforeImages []*types.RecordImage) ([]*types.RecordImage, error) {
  157. if !tm.IsGlobalTx(ctx) {
  158. return nil, nil
  159. }
  160. pc, err := parser.DoParser(execCtx.Query)
  161. if err != nil {
  162. return nil, err
  163. }
  164. if !pc.HasValidStmt() {
  165. return nil, nil
  166. }
  167. execCtx.ParseContext = pc
  168. builder := undo.GetUndologBuilder(pc.ExecutorType)
  169. if builder == nil {
  170. return nil, nil
  171. }
  172. return builder.AfterImage(ctx, execCtx, beforeImages)
  173. }