Cassandra session shared across all storage instances.
| 255 | if cassandra is not None: |
| 256 | |
| 257 | class CassandraSharedSession(object): |
| 258 | """Cassandra session shared across all storage instances.""" |
| 259 | |
| 260 | __session = None |
| 261 | __session_buffer = None |
| 262 | __session_select_buffer = None |
| 263 | |
| 264 | QUERY_CREATE_KEYSPACE = """ |
| 265 | CREATE KEYSPACE IF NOT EXISTS {keyspace} |
| 266 | WITH replication = {replication} |
| 267 | """ |
| 268 | |
| 269 | QUERY_DROP_KEYSPACE = "DROP KEYSPACE IF EXISTS {}" |
| 270 | |
| 271 | @classmethod |
| 272 | def get_session(cls, seeds, **kwargs): |
| 273 | _ = kwargs |
| 274 | keyspace = kwargs["keyspace"] |
| 275 | replication = kwargs["replication"] |
| 276 | |
| 277 | if cls.__session is None: |
| 278 | # Allow dependency injection |
| 279 | session = kwargs.get("session") |
| 280 | if session is None: |
| 281 | cluster = c_cluster.Cluster(seeds) |
| 282 | session = cluster.connect() |
| 283 | cls.__session = session |
| 284 | if cls.__session.keyspace != keyspace: |
| 285 | if kwargs.get("drop_keyspace", False): |
| 286 | cls.__session.execute(cls.QUERY_DROP_KEYSPACE.format(keyspace)) |
| 287 | cls.__session.execute(cls.QUERY_CREATE_KEYSPACE.format( |
| 288 | keyspace=keyspace, |
| 289 | replication=str(replication), |
| 290 | )) |
| 291 | cls.__session.set_keyspace(keyspace) |
| 292 | return cls.__session |
| 293 | |
| 294 | @classmethod |
| 295 | def get_buffer(cls): |
| 296 | if cls.__session_buffer is None: |
| 297 | cls.__session_buffer = [] |
| 298 | return cls.__session_buffer |
| 299 | |
| 300 | @classmethod |
| 301 | def get_select_buffer(cls): |
| 302 | if cls.__session_select_buffer is None: |
| 303 | cls.__session_select_buffer = [] |
| 304 | return cls.__session_select_buffer |
| 305 | |
| 306 | |
| 307 | class CassandraClient(object): |
nothing calls this directly
no outgoing calls
no test coverage detected