Use internal WeakSet if unable to import (users of Python < 2.7)
This commit is contained in:
@@ -10,6 +10,10 @@ from threading import RLock, Thread, Event
|
|||||||
import traceback
|
import traceback
|
||||||
import Queue
|
import Queue
|
||||||
import weakref
|
import weakref
|
||||||
|
try:
|
||||||
|
from weakref import WeakSet
|
||||||
|
except ImportError:
|
||||||
|
from cassandra.util import WeakSet
|
||||||
from functools import partial
|
from functools import partial
|
||||||
from itertools import groupby
|
from itertools import groupby
|
||||||
|
|
||||||
@@ -223,7 +227,7 @@ class Cluster(object):
|
|||||||
# let Session objects be GC'ed (and shutdown) when the user no longer
|
# let Session objects be GC'ed (and shutdown) when the user no longer
|
||||||
# holds a reference. Normally the cycle detector would handle this,
|
# holds a reference. Normally the cycle detector would handle this,
|
||||||
# but implementing __del__ prevents that.
|
# but implementing __del__ prevents that.
|
||||||
self.sessions = weakref.WeakSet()
|
self.sessions = WeakSet()
|
||||||
self.metadata = Metadata(self)
|
self.metadata = Metadata(self)
|
||||||
self.control_connection = None
|
self.control_connection = None
|
||||||
self._prepared_statements = {}
|
self._prepared_statements = {}
|
||||||
|
|||||||
@@ -7,6 +7,10 @@ import time
|
|||||||
from threading import Lock, RLock, Condition
|
from threading import Lock, RLock, Condition
|
||||||
import traceback
|
import traceback
|
||||||
import weakref
|
import weakref
|
||||||
|
try:
|
||||||
|
from weakref import WeakSet
|
||||||
|
except ImportError:
|
||||||
|
from cassandra.util import WeakSet
|
||||||
|
|
||||||
from cassandra import AuthenticationFailed
|
from cassandra import AuthenticationFailed
|
||||||
from cassandra.connection import MAX_STREAM_PER_CONNECTION, ConnectionException
|
from cassandra.connection import MAX_STREAM_PER_CONNECTION, ConnectionException
|
||||||
@@ -211,7 +215,7 @@ class HealthMonitor(object):
|
|||||||
# self._listeners will hold, among other things, references to
|
# self._listeners will hold, among other things, references to
|
||||||
# Cluster objects. To allow those to be GC'ed (and shutdown) even
|
# Cluster objects. To allow those to be GC'ed (and shutdown) even
|
||||||
# though we've implemented __del__, use weak references.
|
# though we've implemented __del__, use weak references.
|
||||||
self._listeners = weakref.WeakSet()
|
self._listeners = WeakSet()
|
||||||
self._lock = RLock()
|
self._lock = RLock()
|
||||||
|
|
||||||
def register(self, listener):
|
def register(self, listener):
|
||||||
|
|||||||
@@ -1,3 +1,7 @@
|
|||||||
|
from __future__ import with_statement
|
||||||
|
|
||||||
|
# OrderedDict from Python 2.7+
|
||||||
|
|
||||||
# Copyright (c) 2009 Raymond Hettinger
|
# Copyright (c) 2009 Raymond Hettinger
|
||||||
#
|
#
|
||||||
# Permission is hereby granted, free of charge, to any person
|
# Permission is hereby granted, free of charge, to any person
|
||||||
@@ -128,3 +132,213 @@ class OrderedDict(dict, DictMixin):
|
|||||||
|
|
||||||
def __ne__(self, other):
|
def __ne__(self, other):
|
||||||
return not self == other
|
return not self == other
|
||||||
|
|
||||||
|
|
||||||
|
# WeakSet from Python 2.7+ (https://code.google.com/p/weakrefset)
|
||||||
|
|
||||||
|
from _weakref import ref
|
||||||
|
|
||||||
|
|
||||||
|
class _IterationGuard(object):
|
||||||
|
# This context manager registers itself in the current iterators of the
|
||||||
|
# weak container, such as to delay all removals until the context manager
|
||||||
|
# exits.
|
||||||
|
# This technique should be relatively thread-safe (since sets are).
|
||||||
|
|
||||||
|
def __init__(self, weakcontainer):
|
||||||
|
# Don't create cycles
|
||||||
|
self.weakcontainer = ref(weakcontainer)
|
||||||
|
|
||||||
|
def __enter__(self):
|
||||||
|
w = self.weakcontainer()
|
||||||
|
if w is not None:
|
||||||
|
w._iterating.add(self)
|
||||||
|
return self
|
||||||
|
|
||||||
|
def __exit__(self, e, t, b):
|
||||||
|
w = self.weakcontainer()
|
||||||
|
if w is not None:
|
||||||
|
s = w._iterating
|
||||||
|
s.remove(self)
|
||||||
|
if not s:
|
||||||
|
w._commit_removals()
|
||||||
|
|
||||||
|
|
||||||
|
class WeakSet(object):
|
||||||
|
def __init__(self, data=None):
|
||||||
|
self.data = set()
|
||||||
|
def _remove(item, selfref=ref(self)):
|
||||||
|
self = selfref()
|
||||||
|
if self is not None:
|
||||||
|
if self._iterating:
|
||||||
|
self._pending_removals.append(item)
|
||||||
|
else:
|
||||||
|
self.data.discard(item)
|
||||||
|
self._remove = _remove
|
||||||
|
# A list of keys to be removed
|
||||||
|
self._pending_removals = []
|
||||||
|
self._iterating = set()
|
||||||
|
if data is not None:
|
||||||
|
self.update(data)
|
||||||
|
|
||||||
|
def _commit_removals(self):
|
||||||
|
l = self._pending_removals
|
||||||
|
discard = self.data.discard
|
||||||
|
while l:
|
||||||
|
discard(l.pop())
|
||||||
|
|
||||||
|
def __iter__(self):
|
||||||
|
with _IterationGuard(self):
|
||||||
|
for itemref in self.data:
|
||||||
|
item = itemref()
|
||||||
|
if item is not None:
|
||||||
|
yield item
|
||||||
|
|
||||||
|
def __len__(self):
|
||||||
|
return sum(x() is not None for x in self.data)
|
||||||
|
|
||||||
|
def __contains__(self, item):
|
||||||
|
return ref(item) in self.data
|
||||||
|
|
||||||
|
def __reduce__(self):
|
||||||
|
return (self.__class__, (list(self),),
|
||||||
|
getattr(self, '__dict__', None))
|
||||||
|
|
||||||
|
__hash__ = None
|
||||||
|
|
||||||
|
def add(self, item):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.add(ref(item, self._remove))
|
||||||
|
|
||||||
|
def clear(self):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.clear()
|
||||||
|
|
||||||
|
def copy(self):
|
||||||
|
return self.__class__(self)
|
||||||
|
|
||||||
|
def pop(self):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
itemref = self.data.pop()
|
||||||
|
except KeyError:
|
||||||
|
raise KeyError('pop from empty WeakSet')
|
||||||
|
item = itemref()
|
||||||
|
if item is not None:
|
||||||
|
return item
|
||||||
|
|
||||||
|
def remove(self, item):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.remove(ref(item))
|
||||||
|
|
||||||
|
def discard(self, item):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.discard(ref(item))
|
||||||
|
|
||||||
|
def update(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
if isinstance(other, self.__class__):
|
||||||
|
self.data.update(other.data)
|
||||||
|
else:
|
||||||
|
for element in other:
|
||||||
|
self.add(element)
|
||||||
|
|
||||||
|
def __ior__(self, other):
|
||||||
|
self.update(other)
|
||||||
|
return self
|
||||||
|
|
||||||
|
# Helper functions for simple delegating methods.
|
||||||
|
def _apply(self, other, method):
|
||||||
|
if not isinstance(other, self.__class__):
|
||||||
|
other = self.__class__(other)
|
||||||
|
newdata = method(other.data)
|
||||||
|
newset = self.__class__()
|
||||||
|
newset.data = newdata
|
||||||
|
return newset
|
||||||
|
|
||||||
|
def difference(self, other):
|
||||||
|
return self._apply(other, self.data.difference)
|
||||||
|
__sub__ = difference
|
||||||
|
|
||||||
|
def difference_update(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
if self is other:
|
||||||
|
self.data.clear()
|
||||||
|
else:
|
||||||
|
self.data.difference_update(ref(item) for item in other)
|
||||||
|
def __isub__(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
if self is other:
|
||||||
|
self.data.clear()
|
||||||
|
else:
|
||||||
|
self.data.difference_update(ref(item) for item in other)
|
||||||
|
return self
|
||||||
|
|
||||||
|
def intersection(self, other):
|
||||||
|
return self._apply(other, self.data.intersection)
|
||||||
|
__and__ = intersection
|
||||||
|
|
||||||
|
def intersection_update(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.intersection_update(ref(item) for item in other)
|
||||||
|
def __iand__(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
self.data.intersection_update(ref(item) for item in other)
|
||||||
|
return self
|
||||||
|
|
||||||
|
def issubset(self, other):
|
||||||
|
return self.data.issubset(ref(item) for item in other)
|
||||||
|
__lt__ = issubset
|
||||||
|
|
||||||
|
def __le__(self, other):
|
||||||
|
return self.data <= set(ref(item) for item in other)
|
||||||
|
|
||||||
|
def issuperset(self, other):
|
||||||
|
return self.data.issuperset(ref(item) for item in other)
|
||||||
|
__gt__ = issuperset
|
||||||
|
|
||||||
|
def __ge__(self, other):
|
||||||
|
return self.data >= set(ref(item) for item in other)
|
||||||
|
|
||||||
|
def __eq__(self, other):
|
||||||
|
if not isinstance(other, self.__class__):
|
||||||
|
return NotImplemented
|
||||||
|
return self.data == set(ref(item) for item in other)
|
||||||
|
|
||||||
|
def symmetric_difference(self, other):
|
||||||
|
return self._apply(other, self.data.symmetric_difference)
|
||||||
|
__xor__ = symmetric_difference
|
||||||
|
|
||||||
|
def symmetric_difference_update(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
if self is other:
|
||||||
|
self.data.clear()
|
||||||
|
else:
|
||||||
|
self.data.symmetric_difference_update(ref(item) for item in other)
|
||||||
|
def __ixor__(self, other):
|
||||||
|
if self._pending_removals:
|
||||||
|
self._commit_removals()
|
||||||
|
if self is other:
|
||||||
|
self.data.clear()
|
||||||
|
else:
|
||||||
|
self.data.symmetric_difference_update(ref(item) for item in other)
|
||||||
|
return self
|
||||||
|
|
||||||
|
def union(self, other):
|
||||||
|
return self._apply(other, self.data.union)
|
||||||
|
__or__ = union
|
||||||
|
|
||||||
|
def isdisjoint(self, other):
|
||||||
|
return len(self.intersection(other)) == 0
|
||||||
Reference in New Issue
Block a user