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

postgresql.go 6.6KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296
  1. package testfixtures
  2. import (
  3. "database/sql"
  4. "fmt"
  5. "strings"
  6. )
  7. type postgreSQL struct {
  8. baseHelper
  9. useAlterConstraint bool
  10. skipResetSequences bool
  11. resetSequencesTo int64
  12. tables []string
  13. sequences []string
  14. nonDeferrableConstraints []pgConstraint
  15. tablesChecksum map[string]string
  16. }
  17. type pgConstraint struct {
  18. tableName string
  19. constraintName string
  20. }
  21. func (h *postgreSQL) init(db *sql.DB) error {
  22. var err error
  23. h.tables, err = h.tableNames(db)
  24. if err != nil {
  25. return err
  26. }
  27. h.sequences, err = h.getSequences(db)
  28. if err != nil {
  29. return err
  30. }
  31. h.nonDeferrableConstraints, err = h.getNonDeferrableConstraints(db)
  32. if err != nil {
  33. return err
  34. }
  35. return nil
  36. }
  37. func (*postgreSQL) paramType() int {
  38. return paramTypeDollar
  39. }
  40. func (*postgreSQL) databaseName(q queryable) (string, error) {
  41. var dbName string
  42. err := q.QueryRow("SELECT current_database()").Scan(&dbName)
  43. return dbName, err
  44. }
  45. func (h *postgreSQL) tableNames(q queryable) ([]string, error) {
  46. var tables []string
  47. sql := `
  48. SELECT pg_namespace.nspname || '.' || pg_class.relname
  49. FROM pg_class
  50. INNER JOIN pg_namespace ON pg_namespace.oid = pg_class.relnamespace
  51. WHERE pg_class.relkind = 'r'
  52. AND pg_namespace.nspname NOT IN ('pg_catalog', 'information_schema')
  53. AND pg_namespace.nspname NOT LIKE 'pg_toast%'
  54. AND pg_namespace.nspname NOT LIKE '\_timescaledb%';
  55. `
  56. rows, err := q.Query(sql)
  57. if err != nil {
  58. return nil, err
  59. }
  60. defer rows.Close()
  61. for rows.Next() {
  62. var table string
  63. if err = rows.Scan(&table); err != nil {
  64. return nil, err
  65. }
  66. tables = append(tables, table)
  67. }
  68. if err = rows.Err(); err != nil {
  69. return nil, err
  70. }
  71. return tables, nil
  72. }
  73. func (h *postgreSQL) getSequences(q queryable) ([]string, error) {
  74. const sql = `
  75. SELECT pg_namespace.nspname || '.' || pg_class.relname AS sequence_name
  76. FROM pg_class
  77. INNER JOIN pg_namespace ON pg_namespace.oid = pg_class.relnamespace
  78. WHERE pg_class.relkind = 'S'
  79. AND pg_namespace.nspname NOT LIKE '\_timescaledb%'
  80. `
  81. rows, err := q.Query(sql)
  82. if err != nil {
  83. return nil, err
  84. }
  85. defer rows.Close()
  86. var sequences []string
  87. for rows.Next() {
  88. var sequence string
  89. if err = rows.Scan(&sequence); err != nil {
  90. return nil, err
  91. }
  92. sequences = append(sequences, sequence)
  93. }
  94. if err = rows.Err(); err != nil {
  95. return nil, err
  96. }
  97. return sequences, nil
  98. }
  99. func (*postgreSQL) getNonDeferrableConstraints(q queryable) ([]pgConstraint, error) {
  100. var constraints []pgConstraint
  101. sql := `
  102. SELECT table_schema || '.' || table_name, constraint_name
  103. FROM information_schema.table_constraints
  104. WHERE constraint_type = 'FOREIGN KEY'
  105. AND is_deferrable = 'NO'
  106. AND table_schema NOT LIKE '\_timescaledb%'
  107. `
  108. rows, err := q.Query(sql)
  109. if err != nil {
  110. return nil, err
  111. }
  112. defer rows.Close()
  113. for rows.Next() {
  114. var constraint pgConstraint
  115. if err = rows.Scan(&constraint.tableName, &constraint.constraintName); err != nil {
  116. return nil, err
  117. }
  118. constraints = append(constraints, constraint)
  119. }
  120. if err = rows.Err(); err != nil {
  121. return nil, err
  122. }
  123. return constraints, nil
  124. }
  125. func (h *postgreSQL) disableTriggers(db *sql.DB, loadFn loadFunction) (err error) {
  126. defer func() {
  127. // re-enable triggers after load
  128. var sql string
  129. for _, table := range h.tables {
  130. sql += fmt.Sprintf("ALTER TABLE %s ENABLE TRIGGER ALL;", h.quoteKeyword(table))
  131. }
  132. if _, err2 := db.Exec(sql); err2 != nil && err == nil {
  133. err = err2
  134. }
  135. }()
  136. tx, err := db.Begin()
  137. if err != nil {
  138. return err
  139. }
  140. var sql string
  141. for _, table := range h.tables {
  142. sql += fmt.Sprintf("ALTER TABLE %s DISABLE TRIGGER ALL;", h.quoteKeyword(table))
  143. }
  144. if _, err = tx.Exec(sql); err != nil {
  145. return err
  146. }
  147. if err = loadFn(tx); err != nil {
  148. tx.Rollback()
  149. return err
  150. }
  151. return tx.Commit()
  152. }
  153. func (h *postgreSQL) makeConstraintsDeferrable(db *sql.DB, loadFn loadFunction) (err error) {
  154. defer func() {
  155. // ensure constraint being not deferrable again after load
  156. var sql string
  157. for _, constraint := range h.nonDeferrableConstraints {
  158. sql += fmt.Sprintf("ALTER TABLE %s ALTER CONSTRAINT %s NOT DEFERRABLE;", h.quoteKeyword(constraint.tableName), h.quoteKeyword(constraint.constraintName))
  159. }
  160. if _, err2 := db.Exec(sql); err2 != nil && err == nil {
  161. err = err2
  162. }
  163. }()
  164. var sql string
  165. for _, constraint := range h.nonDeferrableConstraints {
  166. sql += fmt.Sprintf("ALTER TABLE %s ALTER CONSTRAINT %s DEFERRABLE;", h.quoteKeyword(constraint.tableName), h.quoteKeyword(constraint.constraintName))
  167. }
  168. if _, err := db.Exec(sql); err != nil {
  169. return err
  170. }
  171. tx, err := db.Begin()
  172. if err != nil {
  173. return err
  174. }
  175. defer tx.Rollback()
  176. if _, err = tx.Exec("SET CONSTRAINTS ALL DEFERRED"); err != nil {
  177. return err
  178. }
  179. if err = loadFn(tx); err != nil {
  180. return err
  181. }
  182. return tx.Commit()
  183. }
  184. func (h *postgreSQL) disableReferentialIntegrity(db *sql.DB, loadFn loadFunction) (err error) {
  185. // ensure sequences being reset after load
  186. if !h.skipResetSequences {
  187. defer func() {
  188. if err2 := h.resetSequences(db); err2 != nil && err == nil {
  189. err = err2
  190. }
  191. }()
  192. }
  193. if h.useAlterConstraint {
  194. return h.makeConstraintsDeferrable(db, loadFn)
  195. }
  196. return h.disableTriggers(db, loadFn)
  197. }
  198. func (h *postgreSQL) resetSequences(db *sql.DB) error {
  199. resetSequencesTo := h.resetSequencesTo
  200. if resetSequencesTo == 0 {
  201. resetSequencesTo = 10000
  202. }
  203. for _, sequence := range h.sequences {
  204. _, err := db.Exec(fmt.Sprintf("SELECT SETVAL('%s', %d)", sequence, resetSequencesTo))
  205. if err != nil {
  206. return err
  207. }
  208. }
  209. return nil
  210. }
  211. func (h *postgreSQL) isTableModified(q queryable, tableName string) (bool, error) {
  212. checksum, err := h.getChecksum(q, tableName)
  213. if err != nil {
  214. return false, err
  215. }
  216. oldChecksum := h.tablesChecksum[tableName]
  217. return oldChecksum == "" || checksum != oldChecksum, nil
  218. }
  219. func (h *postgreSQL) afterLoad(q queryable) error {
  220. if h.tablesChecksum != nil {
  221. return nil
  222. }
  223. h.tablesChecksum = make(map[string]string, len(h.tables))
  224. for _, t := range h.tables {
  225. checksum, err := h.getChecksum(q, t)
  226. if err != nil {
  227. return err
  228. }
  229. h.tablesChecksum[t] = checksum
  230. }
  231. return nil
  232. }
  233. func (h *postgreSQL) getChecksum(q queryable, tableName string) (string, error) {
  234. sqlStr := fmt.Sprintf(`
  235. SELECT md5(CAST((array_agg(t.*)) AS TEXT))
  236. FROM %s AS t
  237. `,
  238. h.quoteKeyword(tableName),
  239. )
  240. var checksum sql.NullString
  241. if err := q.QueryRow(sqlStr).Scan(&checksum); err != nil {
  242. return "", err
  243. }
  244. return checksum.String, nil
  245. }
  246. func (*postgreSQL) quoteKeyword(s string) string {
  247. parts := strings.Split(s, ".")
  248. for i, p := range parts {
  249. parts[i] = fmt.Sprintf(`"%s"`, p)
  250. }
  251. return strings.Join(parts, ".")
  252. }