97
98
// DatabaseName from PersistenceTestCluster interface
99
>
func (s *TestCluster) DatabaseName() string {
test.go
100
>
return s.keyspace
101
>
}
102
103
// SetupTestDatabase from PersistenceTestCluster interface
104
>
func (s *TestCluster) SetupTestDatabase() {
test.go
105
>
s.CreateSession("system")
106
>
s.CreateDatabase()
107
>
s.CreateSession(s.DatabaseName())
108
>
schemaDir := s.schemaDir + "/"
109
>
110
>
if !strings.HasPrefix(schemaDir, "/") && !strings.HasPrefix(schemaDir, "../") {
111
>
temporalPackageDir := testutils.GetRepoRootDirectory()
112
>
schemaDir = path.Join(temporalPackageDir, schemaDir)
113
>
}
114
115
>
s.LoadSchema(path.Join(schemaDir, "temporal", "schema.cql"))
test.go
116
>
s.loadSchemaVersion()
117
}
118
119
// TearDownTestDatabase from PersistenceTestCluster interface
120
>
func (s *TestCluster) TearDownTestDatabase() {
test.go
121
>
s.DropDatabase()
122
>
s.session.Close()
123
>
}
124
125
// CreateSession from PersistenceTestCluster interface
126
func (s *TestCluster) CreateSession(
127
keyspace string,
129
>
if s.session != nil {
130
>
s.session.Close()
131
>
}
132
134
>
op := func() error {
135
>
session, err := commongocql.NewSession(
136
>
func() (*gocql.ClusterConfig, error) {
137
>
return commongocql.NewCassandraCluster(
138
>
config.Cassandra{
139
>
Hosts: s.cfg.Hosts,
140
>
Port: s.cfg.Port,
141
>
User: s.cfg.User,
142
>
Password: s.cfg.Password,
143
>
Keyspace: keyspace,
144
>
Consistency: &config.CassandraStoreConsistency{
145
>
Default: &config.CassandraConsistencySettings{
146
>
Consistency: "ONE",
147
>
},
148
>
},
149
>
ConnectTimeout: s.cfg.ConnectTimeout,
150
>
},
151
>
resolver.NewNoopResolver(),
152
>
)
153
>
},
154
log.NewNoopLogger(),
155
metrics.NoopMetricsHandler,
156
)
158
>
s.session = session
159
>
}
160
>
return err
161
}
162
>
err = backoff.ThrottleRetry(
test.go
163
>
op,
164
>
backoff.NewExponentialRetryPolicy(time.Second).WithExpirationInterval(time.Minute),
165
>
nil,
166
>
)
167
>
if err != nil {
168
s.logger.Fatal("CreateSession", tag.Error(err))
169
}
170
>
s.logger.Debug("created session", tag.String("keyspace", keyspace))
test.go
171
}
172
173
// CreateDatabase from PersistenceTestCluster interface
174
>
func (s *TestCluster) CreateDatabase() {
test.go
175
>
err := CreateCassandraKeyspace(s.session, s.DatabaseName(), 1, true, s.logger)
176
>
if err != nil {
177
s.logger.Fatal("CreateCassandraKeyspace", tag.Error(err))
178
}
179
>
s.logger.Info("created database", tag.String("database", s.DatabaseName()))
test.go
180
}
181
182
// DropDatabase from PersistenceTestCluster interface
183
>
func (s *TestCluster) DropDatabase() {
test.go
184
>
err := DropCassandraKeyspace(s.session, s.DatabaseName(), s.logger)
185
>
if err != nil && !strings.Contains(err.Error(), "AlreadyExists") {
186
s.logger.Fatal("DropCassandraKeyspace", tag.Error(err))
187
}
188
>
s.logger.Info("dropped database", tag.String("database", s.DatabaseName()))
test.go
189
}
190
191
// LoadSchema from PersistenceTestCluster interface
192
>
func (s *TestCluster) LoadSchema(schemaFile string) {
test.go
193
>
statements, err := p.LoadAndSplitQuery([]string{schemaFile})
194
>
if err != nil {
195
s.logger.Fatal("LoadSchema", tag.Error(err))
196
}
197
>
for _, stmt := range statements {
test.go
198
>
if err = s.session.Query(stmt).Exec(); err != nil {
199
s.logger.Fatal("LoadSchema", tag.Error(err))
200
}
201
}
202
>
s.logger.Info("loaded schema")
test.go
203
}
204
205
>
func (s *TestCluster) loadSchemaVersion() {
test.go
206
>
s.createSchemaVersionTables()
207
>
s.updateSchemaVersion(cassandraschema.Version, cassandraschema.Version)
208
>
s.writeSchemaUpdateLog("0", cassandraschema.Version, "", "initial version")
209
>
s.logger.Info("loaded schema version", tag.String("version", cassandraschema.Version))
210
>
}
211
212
>
func (s *TestCluster) createSchemaVersionTables() {
test.go
213
>
s.execSchemaVersionQuery(createSchemaVersionTableCQL)
214
>
s.execSchemaVersionQuery(createSchemaUpdateHistoryTableCQL)
215
>
}
216
217
>
func (s *TestCluster) updateSchemaVersion(newVersion string, minCompatibleVersion string) {
test.go
218
>
now := time.Now().UTC()
219
>
s.execSchemaVersionQuery(
220
>
writeSchemaVersionCQL,
221
>
s.keyspace, now, newVersion, minCompatibleVersion)
222
>
}
223
224
>
func (s *TestCluster) writeSchemaUpdateLog(oldVersion string, newVersion string, manifestMD5 string, description string) {
test.go
225
>
now := time.Now().UTC()
226
>
s.execSchemaVersionQuery(
227
>
writeSchemaUpdateHistoryCQL,
228
>
now.Year(), int(now.Month()), now, oldVersion, newVersion, manifestMD5, description)
229
>
}
230
231
>
func (s *TestCluster) execSchemaVersionQuery(stmt string, args ...any) {
test.go
232
>
if err := s.session.Query(stmt, args...).Exec(); err != nil {
233
s.logger.Fatal("loadSchemaVersion", tag.Error(err))
234
}