queue_store.go ×10

Frontier kind: Code frontier

unlabeled · c_9a9835ae56ed

38 tests · 4513 LOC · 198 files · introduces 0 tests · 127 LOC · 13 files

Introduces — evidence that enters the hierarchy at this concept

Code
27 ranges127 lines · 13 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
834 ranges4513 lines · 198 files · Browse complete extent
All tests (intent)
38 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

13 files ranked by introduced lines: 127 introduced LOC across 27 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/cassandra/queue_store.go 51 introduced LOC · 10 ranges

Open complete file

41 session gocql.Session,
42 logger log.Logger,
43 > ) (persistence.Queue, error) { queue_store.go
44 > return &QueueStore{
45 > queueType: queueType,
46 > session: session,
47 > logger: logger,
48 > serializer: serialization.NewSerializer(),
49 > }, nil
50 > }
51
52 func (q *QueueStore) Init(
53 ctx context.Context,
54 blob *commonpb.DataBlob,
55 > ) error { queue_store.go
56 > if err := q.initializeQueueMetadata(ctx, blob); err != nil {
57 return err
58 }
59 > return q.initializeDLQMetadata(ctx, blob) queue_store.go
60 }
61
270 queueType persistence.QueueType,
271 blob *commonpb.DataBlob,
272 > ) error { queue_store.go
273 >
274 > version := 0
275 > // TODO: remove once cluster_ack_level is removed from DB
276 > clusterAckLevels := map[string]int64{}
277 > query := q.session.Query(templateInsertQueueMetadataQuery,
278 > queueType,
279 > clusterAckLevels,
280 > blob.Data,
281 > blob.EncodingType.String(),
282 > version,
283 > ).WithContext(ctx)
284 > _, err := query.MapScanCAS(make(map[string]any))
285 > if err != nil {
286 return fmt.Errorf("failed to insert initial queue metadata record: %v, Type: %v", err, queueType)
287 }
288 // it's ok if the query is not applied, which means that the record exists already.
289 > return nil queue_store.go
290 }
291
293 ctx context.Context,
294 queueType persistence.QueueType,
295 > ) (*persistence.InternalQueueMetadata, error) { queue_store.go
296 >
297 > query := q.session.Query(templateGetQueueMetadataQuery, queueType).WithContext(ctx)
298 > message := make(map[string]any)
299 > err := query.MapScan(message)
300 > if err != nil {
301 > return nil, err
302 > }
303
304 return convertQueueMetadata(message, q.serializer)
336 }
337
338 > func (q *QueueStore) Close() { queue_store.go
339 > if q.session != nil {
340 > q.session.Close()
341 > }
342 }
343
344 > func (q *QueueStore) getDLQTypeFromQueueType() persistence.QueueType { queue_store.go
345 > return -q.queueType
346 > }
347
348 func (q *QueueStore) initializeQueueMetadata(
349 ctx context.Context,
350 blob *commonpb.DataBlob,
351 > ) error { queue_store.go
352 > _, err := q.getQueueMetadata(ctx, q.queueType)
353 > if gocql.IsNotFoundError(err) {
354 > return q.insertInitialQueueMetadataRecord(ctx, q.queueType, blob)
355 > }
356 return err
357 }
360 ctx context.Context,
361 blob *commonpb.DataBlob,
362 > ) error { queue_store.go
363 > _, err := q.getQueueMetadata(ctx, q.getDLQTypeFromQueueType())
364 > if gocql.IsNotFoundError(err) {
365 > return q.insertInitialQueueMetadataRecord(ctx, q.getDLQTypeFromQueueType(), blob)
366 > }
367 return err
368 }
go.temporal.io/server/common/persistence/cassandra/factory.go 12 introduced LOC · 4 ranges

Open complete file

89
90 // NewMetadataStore returns a metadata store
91 > func (f *Factory) NewMetadataStore() (p.MetadataStore, error) { factory.go
92 > return NewMetadataStore(f.clusterName, f.session, f.logger)
93 > }
94
95 // NewClusterMetadataStore returns a metadata store
96 > func (f *Factory) NewClusterMetadataStore() (p.ClusterMetadataStore, error) { factory.go
97 > return NewClusterMetadataStore(f.session, f.logger)
98 > }
99
100 // NewExecutionStore returns a new ExecutionStore.
104
105 // NewQueue returns a new queue backed by cassandra
106 > func (f *Factory) NewQueue(queueType p.QueueType) (p.Queue, error) { factory.go
107 > return NewQueueStore(queueType, f.session, f.logger)
108 > }
109
110 // NewQueueV2 returns a new data-access object for queues and messages stored in Cassandra. It will never return an
115
116 // NewNexusEndpointStore returns a new NexusEndpointStore
117 > func (f *Factory) NewNexusEndpointStore() (p.NexusEndpointStore, error) { factory.go
118 > return NewNexusEndpointStore(f.session, f.logger), nil
119 > }
120
121 // Close closes the factory
go.temporal.io/server/common/persistence/cassandra/metadata_store.go 11 introduced LOC · 2 ranges

Open complete file

84 session gocql.Session,
85 logger log.Logger,
86 > ) (p.MetadataStore, error) { metadata_store.go
87 > return &MetadataStore{
88 > currentClusterName: currentClusterName,
89 > session: session,
90 > logger: logger,
91 > }, nil
92 > }
93
94 // CreateNamespace create a namespace
504 }
505
506 > func (m *MetadataStore) Close() { metadata_store.go
507 > if m.session != nil {
508 > m.session.Close()
509 > }
510 }
511
go.temporal.io/server/common/persistence/cassandra/cluster_metadata_store.go 10 introduced LOC · 2 ranges

Open complete file

54 session gocql.Session,
55 logger log.Logger,
56 > ) (p.ClusterMetadataStore, error) { cluster_metadata_store.go
57 > return &ClusterMetadataStore{
58 > session: session,
59 > logger: logger,
60 > }, nil
61 > }
62
63 func (m *ClusterMetadataStore) ListClusterMetadata(
279 }
280
281 > func (m *ClusterMetadataStore) Close() { cluster_metadata_store.go
282 > if m.session != nil {
283 > m.session.Close()
284 > }
285 }
go.temporal.io/server/common/persistence/cassandra/execution_store.go 10 introduced LOC · 1 range

Open complete file

153 }
154
155 > func (d *ExecutionStore) Close() { execution_store.go
156 > if d.HistoryStore.Session != nil {
157 > d.HistoryStore.Session.Close()
158 > }
159 > if d.MutableStateStore.Session != nil {
160 > d.MutableStateStore.Session.Close()
161 > }
162 > if d.MutableStateTaskStore.Session != nil {
163 > d.MutableStateTaskStore.Session.Close()
164 > }
165 }
go.temporal.io/server/common/persistence/cassandra/test.go 10 introduced LOC · 1 range

Open complete file

85
86 // Config returns the persistence config for connecting to this test cluster
87 > func (s *TestCluster) Config() config.Persistence { test.go
88 > cfg := s.cfg
89 > return config.Persistence{
90 > DefaultStore: "test",
91 > DataStores: map[string]config.DataStore{
92 > "test": {Cassandra: &cfg, FaultInjection: s.faultInjection},
93 > },
94 > TransactionSizeLimit: dynamicconfig.GetIntPropertyFn(primitives.DefaultTransactionSizeLimit),
95 > }
96 > }
97
98 // DatabaseName from PersistenceTestCluster interface
go.temporal.io/server/common/persistence/persistence-tests/persistence_test_base.go 5 introduced LOC · 1 range

Open complete file

122
123 // NewTestBaseWithCassandra returns a persistence test base backed by cassandra datastore
124 > func NewTestBaseWithCassandra(options *TestBaseOptions) *TestBase { persistence_test_base.go
125 > logger := log.NewTestLogger()
126 > testCluster := NewTestClusterForCassandra(options, logger)
127 > return NewTestBaseForCluster(testCluster, logger)
128 > }
129
130 func NewTestClusterForCassandra(options *TestBaseOptions, logger log.Logger) *cassandra.TestCluster {
go.temporal.io/server/common/persistence/cassandra/matching_task_store_v1.go 4 introduced LOC · 1 range

Open complete file

225 }
226
227 > func (d *matchingTaskStoreV1) Close() { matching_task_store_v1.go
228 > if d.Session != nil {
229 > d.Session.Close()
230 > }
231 }
go.temporal.io/server/common/persistence/cassandra/nexus_endpoint_store.go 4 introduced LOC · 1 range

Open complete file

57 }
58
59 > func (s *NexusEndpointStore) Close() { nexus_endpoint_store.go
60 > if s.session != nil {
61 > s.session.Close()
62 > }
63 }
64
go.temporal.io/server/common/persistence/cassandra/shard_store.go 4 introduced LOC · 1 range

Open complete file

172 }
173
174 > func (d *ShardStore) Close() { shard_store.go
175 > if d.Session != nil {
176 > d.Session.Close()
177 > }
178 }
go.temporal.io/server/common/persistence/client/fx.go 2 introduced LOC · 1 range

Open complete file

198 defaultStoreCfg := cfg.DataStores[cfg.DefaultStore]
199 switch {
200 > case defaultStoreCfg.Cassandra != nil: fx.go
201 > dataStoreFactory = cassandra.NewFactory(*defaultStoreCfg.Cassandra, r, string(clusterName), logger, metricsHandler, serializer)
202 case defaultStoreCfg.SQL != nil:
203 dataStoreFactory = sql.NewFactory(*defaultStoreCfg.SQL, r, string(clusterName), logger, metricsHandler, serializer)
go.temporal.io/server/common/persistence/nosql/nosqlplugin/cassandra/gocql/client.go 2 introduced LOC · 1 range

Open complete file

126
127 if cfg.MaxConns > 0 {
128 > cluster.NumConns = cfg.MaxConns client.go
129 > }
130
131 cluster.ConnectTimeout = 10 * time.Second * debug.TimeoutMultiplier
go.temporal.io/server/common/persistence/nosql/nosqlplugin/cassandra/gocql/session.go 2 introduced LOC · 1 range

Open complete file

169 common.DaemonStatusStopped,
170 ) {
171 > return session.go
172 > }
173 s.Value.Load().(*gocql.Session).Close()
174 }