tested django-newsletter
This commit is contained in:
@@ -0,0 +1,37 @@
|
||||
############################################################################
|
||||
#
|
||||
# Copyright (c) 2001, 2002, 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
############################################################################
|
||||
"""Exported transaction functions.
|
||||
|
||||
$Id$
|
||||
"""
|
||||
|
||||
from transaction._transaction import Transaction
|
||||
from transaction._manager import TransactionManager
|
||||
from transaction._manager import ThreadTransactionManager
|
||||
|
||||
# NB: "with transaction:" does not work under Python 3 because they worked
|
||||
# really hard to break looking up special methods like __enter__ and __exit__
|
||||
# via getattr and getattribute; see http://bugs.python.org/issue12022. On
|
||||
# Python 3, you must use ``with transaction.manager`` instead.
|
||||
|
||||
manager = ThreadTransactionManager()
|
||||
get = __enter__ = manager.get
|
||||
begin = manager.begin
|
||||
commit = manager.commit
|
||||
abort = manager.abort
|
||||
__exit__ = manager.__exit__
|
||||
doom = manager.doom
|
||||
isDoomed = manager.isDoomed
|
||||
savepoint = manager.savepoint
|
||||
attempts = manager.attempts
|
||||
@@ -0,0 +1,74 @@
|
||||
import sys
|
||||
|
||||
|
||||
PY3 = sys.version_info[0] == 3
|
||||
JYTHON = sys.platform.startswith('java')
|
||||
|
||||
if PY3:
|
||||
text_type = str
|
||||
else: # pragma: no cover
|
||||
# py2
|
||||
text_type = unicode
|
||||
|
||||
def bytes_(s, encoding='latin-1', errors='strict'):
|
||||
if isinstance(s, text_type):
|
||||
s = s.encode(encoding, errors)
|
||||
return s
|
||||
|
||||
def text_(s):
|
||||
if not isinstance(s, text_type):
|
||||
s = s.decode('utf-8')
|
||||
return s
|
||||
|
||||
if PY3:
|
||||
def native_(s, encoding='latin-1', errors='strict'):
|
||||
if isinstance(s, text_type):
|
||||
return s
|
||||
return str(s, encoding, errors)
|
||||
else: # pragma: no cover
|
||||
def native_(s, encoding='latin-1', errors='strict'):
|
||||
if isinstance(s, text_type):
|
||||
return s.encode(encoding, errors)
|
||||
return str(s)
|
||||
|
||||
if PY3:
|
||||
from io import StringIO
|
||||
else: # pragma: no cover
|
||||
from io import BytesIO
|
||||
# Prevent crashes in IPython when writing tracebacks if a commit fails
|
||||
# ref: https://github.com/ipython/ipython/issues/9126#issuecomment-174966638
|
||||
class StringIO(BytesIO):
|
||||
def write(self, s):
|
||||
s = native_(s, encoding='utf-8')
|
||||
super(StringIO, self).write(s)
|
||||
|
||||
|
||||
if PY3:
|
||||
def reraise(tp, value, tb=None):
|
||||
if value.__traceback__ is not tb: # pragma: no cover
|
||||
raise value.with_traceback(tb)
|
||||
raise value
|
||||
|
||||
else: # pragma: no cover
|
||||
def exec_(code, globs=None, locs=None):
|
||||
"""Execute code in a namespace."""
|
||||
if globs is None:
|
||||
frame = sys._getframe(1)
|
||||
globs = frame.f_globals
|
||||
if locs is None:
|
||||
locs = frame.f_locals
|
||||
del frame
|
||||
elif locs is None:
|
||||
locs = globs
|
||||
exec("""exec code in globs, locs""")
|
||||
|
||||
exec_("""def reraise(tp, value, tb=None):
|
||||
raise tp, value, tb
|
||||
""")
|
||||
|
||||
|
||||
try:
|
||||
from threading import get_ident as get_thread_ident
|
||||
except ImportError: # pragma: no cover
|
||||
# py2
|
||||
from thread import get_ident as get_thread_ident
|
||||
@@ -0,0 +1,313 @@
|
||||
############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
############################################################################
|
||||
"""A TransactionManager controls transaction boundaries.
|
||||
|
||||
It coordinates application code and resource managers, so that they
|
||||
are associated with the right transaction.
|
||||
"""
|
||||
import sys
|
||||
import threading
|
||||
|
||||
from zope.interface import implementer
|
||||
|
||||
from transaction.interfaces import AlreadyInTransaction
|
||||
from transaction.interfaces import ITransactionManager
|
||||
from transaction.interfaces import NoTransaction
|
||||
from transaction.interfaces import TransientError
|
||||
from transaction.weakset import WeakSet
|
||||
from transaction._compat import reraise
|
||||
from transaction._compat import text_
|
||||
from transaction._transaction import Transaction
|
||||
|
||||
|
||||
# We have to remember sets of synch objects, especially Connections.
|
||||
# But we don't want mere registration with a transaction manager to
|
||||
# keep a synch object alive forever; in particular, it's common
|
||||
# practice not to explicitly close Connection objects, and keeping
|
||||
# a Connection alive keeps a potentially huge number of other objects
|
||||
# alive (e.g., the cache, and everything reachable from it too).
|
||||
# Therefore we use "weak sets" internally.
|
||||
|
||||
# Call the ISynchronizer newTransaction() method on every element of
|
||||
# WeakSet synchs.
|
||||
# A transaction manager needs to do this whenever begin() is called.
|
||||
# Since it would be good if tm.get() returned the new transaction while
|
||||
# newTransaction() is running, calling this has to be delayed until after
|
||||
# the transaction manager has done whatever it needs to do to make its
|
||||
# get() return the new txn.
|
||||
def _new_transaction(txn, synchs):
|
||||
if synchs:
|
||||
synchs.map(lambda s: s.newTransaction(txn))
|
||||
|
||||
# Important: we must always pass a WeakSet (even if empty) to the Transaction
|
||||
# constructor: synchronizers are registered with the TM, but the
|
||||
# ISynchronizer xyzCompletion() methods are called by Transactions without
|
||||
# consulting the TM, so we need to pass a mutable collection of synchronizers
|
||||
# so that Transactions "see" synchronizers that get registered after the
|
||||
# Transaction object is constructed.
|
||||
|
||||
|
||||
@implementer(ITransactionManager)
|
||||
class TransactionManager(object):
|
||||
|
||||
def __init__(self, explicit=False):
|
||||
self.explicit = explicit
|
||||
self._txn = None
|
||||
self._synchs = WeakSet()
|
||||
|
||||
def begin(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
if self._txn is not None:
|
||||
if self.explicit:
|
||||
raise AlreadyInTransaction()
|
||||
self._txn.abort()
|
||||
txn = self._txn = Transaction(self._synchs, self)
|
||||
_new_transaction(txn, self._synchs)
|
||||
return txn
|
||||
|
||||
__enter__ = lambda self: self.begin()
|
||||
|
||||
def get(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
if self._txn is None:
|
||||
if self.explicit:
|
||||
raise NoTransaction()
|
||||
self._txn = Transaction(self._synchs, self)
|
||||
return self._txn
|
||||
|
||||
def free(self, txn):
|
||||
if txn is not self._txn:
|
||||
raise ValueError("Foreign transaction")
|
||||
self._txn = None
|
||||
|
||||
def registerSynch(self, synch):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
self._synchs.add(synch)
|
||||
if self._txn is not None:
|
||||
synch.newTransaction(self._txn)
|
||||
|
||||
def unregisterSynch(self, synch):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
self._synchs.remove(synch)
|
||||
|
||||
def clearSynchs(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
self._synchs.clear()
|
||||
|
||||
def registeredSynchs(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return bool(self._synchs)
|
||||
|
||||
def isDoomed(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return self.get().isDoomed()
|
||||
|
||||
def doom(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return self.get().doom()
|
||||
|
||||
def commit(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return self.get().commit()
|
||||
|
||||
def abort(self):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return self.get().abort()
|
||||
|
||||
def __exit__(self, t, v, tb):
|
||||
if v is None:
|
||||
self.commit()
|
||||
else:
|
||||
self.abort()
|
||||
|
||||
def savepoint(self, optimistic=False):
|
||||
""" See ITransactionManager.
|
||||
"""
|
||||
return self.get().savepoint(optimistic)
|
||||
|
||||
def attempts(self, number=3):
|
||||
if number <= 0:
|
||||
raise ValueError("number must be positive")
|
||||
while number:
|
||||
number -= 1
|
||||
if number:
|
||||
attempt = Attempt(self)
|
||||
yield attempt
|
||||
if attempt.success:
|
||||
break
|
||||
else:
|
||||
yield self
|
||||
|
||||
def _retryable(self, error_type, error):
|
||||
if issubclass(error_type, TransientError):
|
||||
return True
|
||||
|
||||
for dm in self.get()._resources:
|
||||
should_retry = getattr(dm, 'should_retry', None)
|
||||
if (should_retry is not None) and should_retry(error):
|
||||
return True
|
||||
return False
|
||||
|
||||
run_no_func_types = int, type(None)
|
||||
def run(self, func=None, tries=3):
|
||||
if isinstance(func, self.run_no_func_types):
|
||||
if func is not None:
|
||||
tries = func
|
||||
return lambda func: self.run(func, tries)
|
||||
|
||||
if tries <= 0:
|
||||
raise ValueError("tries must be positive")
|
||||
|
||||
# These are ordinarily native strings, but that's
|
||||
# not required. A callable class could override them
|
||||
# to anything, and a Python 2.7 file could have
|
||||
# imported `from __future__ import unicode_literals`
|
||||
# which gets unicode docstrings.
|
||||
name = func.__name__
|
||||
doc = func.__doc__
|
||||
|
||||
name = text_(name) if name else u''
|
||||
doc = text_(doc) if doc else u''
|
||||
|
||||
if name != u'_':
|
||||
if doc:
|
||||
doc = name + u'\n\n' + doc
|
||||
else:
|
||||
doc = name
|
||||
|
||||
for i in range(1, tries + 1): # pragma: no branch
|
||||
txn = self.begin()
|
||||
if doc:
|
||||
txn.note(doc)
|
||||
|
||||
try:
|
||||
result = func()
|
||||
txn.commit()
|
||||
except Exception as v:
|
||||
if i == tries:
|
||||
raise # that was our last chance
|
||||
retry = self._retryable(v.__class__, v)
|
||||
txn.abort()
|
||||
if not retry:
|
||||
raise
|
||||
else:
|
||||
return result
|
||||
|
||||
|
||||
@implementer(ITransactionManager)
|
||||
class ThreadTransactionManager(threading.local):
|
||||
"""
|
||||
Thread-local transaction manager.
|
||||
|
||||
A thread-local transaction manager can be used as a global
|
||||
variable, but has a separate copy for each thread.
|
||||
|
||||
Advanced applications can use the `manager` attribute to get a
|
||||
wrapped TransactionManager to allow cross-thread calls for
|
||||
graceful shutdown of data managers.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.manager = TransactionManager()
|
||||
|
||||
@property
|
||||
def explicit(self):
|
||||
return self.manager.explicit
|
||||
|
||||
@explicit.setter
|
||||
def explicit(self, v):
|
||||
self.manager.explicit = v
|
||||
|
||||
def begin(self):
|
||||
return self.manager.begin()
|
||||
|
||||
def get(self):
|
||||
return self.manager.get()
|
||||
|
||||
def __enter__(self):
|
||||
return self.manager.__enter__()
|
||||
|
||||
def commit(self):
|
||||
return self.manager.commit()
|
||||
|
||||
def abort(self):
|
||||
return self.manager.abort()
|
||||
|
||||
def __exit__(self, t, v, tb):
|
||||
return self.manager.__exit__(t, v, tb)
|
||||
|
||||
def doom(self):
|
||||
return self.manager.doom()
|
||||
|
||||
def isDoomed(self):
|
||||
return self.manager.isDoomed()
|
||||
|
||||
def savepoint(self, optimistic=False):
|
||||
return self.manager.savepoint(optimistic)
|
||||
|
||||
def registerSynch(self, synch):
|
||||
return self.manager.registerSynch(synch)
|
||||
|
||||
def unregisterSynch(self, synch):
|
||||
return self.manager.unregisterSynch(synch)
|
||||
|
||||
def clearSynchs(self):
|
||||
return self.manager.clearSynchs()
|
||||
|
||||
def registeredSynchs(self):
|
||||
return self.manager.registeredSynchs()
|
||||
|
||||
def attempts(self, number=3):
|
||||
return self.manager.attempts(number)
|
||||
|
||||
def run(self, func=None, tries=3):
|
||||
return self.manager.run(func, tries)
|
||||
|
||||
class Attempt(object):
|
||||
|
||||
success = False
|
||||
|
||||
def __init__(self, manager):
|
||||
self.manager = manager
|
||||
|
||||
def _retry_or_raise(self, t, v, tb):
|
||||
retry = self.manager._retryable(t, v)
|
||||
self.manager.abort()
|
||||
if retry:
|
||||
return retry # suppress the exception if necessary
|
||||
reraise(t, v, tb) # otherwise reraise the exception
|
||||
|
||||
def __enter__(self):
|
||||
return self.manager.__enter__()
|
||||
|
||||
def __exit__(self, t, v, tb):
|
||||
if v is None:
|
||||
try:
|
||||
self.manager.commit()
|
||||
except:
|
||||
return self._retry_or_raise(*sys.exc_info())
|
||||
else:
|
||||
self.success = True
|
||||
else:
|
||||
return self._retry_or_raise(t, v, tb)
|
||||
@@ -0,0 +1,769 @@
|
||||
############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
############################################################################
|
||||
import binascii
|
||||
import logging
|
||||
import sys
|
||||
import warnings
|
||||
import weakref
|
||||
import traceback
|
||||
|
||||
from zope.interface import implementer
|
||||
|
||||
from transaction.weakset import WeakSet
|
||||
from transaction.interfaces import TransactionFailedError
|
||||
from transaction import interfaces
|
||||
from transaction._compat import reraise
|
||||
from transaction._compat import get_thread_ident
|
||||
from transaction._compat import native_
|
||||
from transaction._compat import bytes_
|
||||
from transaction._compat import StringIO
|
||||
from transaction._compat import text_type
|
||||
|
||||
_marker = object()
|
||||
|
||||
_TB_BUFFER = None #unittests may hook
|
||||
def _makeTracebackBuffer(): #pragma NO COVER
|
||||
if _TB_BUFFER is not None:
|
||||
return _TB_BUFFER
|
||||
return StringIO()
|
||||
|
||||
_LOGGER = None #unittests may hook
|
||||
def _makeLogger(): #pragma NO COVER
|
||||
if _LOGGER is not None:
|
||||
return _LOGGER
|
||||
return logging.getLogger("txn.%d" % get_thread_ident())
|
||||
|
||||
|
||||
# The point of this is to avoid hiding exceptions (which the builtin
|
||||
# hasattr() does).
|
||||
def myhasattr(obj, attr):
|
||||
return getattr(obj, attr, _marker) is not _marker
|
||||
|
||||
class Status:
|
||||
# ACTIVE is the initial state.
|
||||
ACTIVE = "Active"
|
||||
|
||||
COMMITTING = "Committing"
|
||||
COMMITTED = "Committed"
|
||||
|
||||
DOOMED = "Doomed"
|
||||
|
||||
# commit() or commit(True) raised an exception. All further attempts
|
||||
# to commit or join this transaction will raise TransactionFailedError.
|
||||
COMMITFAILED = "Commit failed"
|
||||
|
||||
@implementer(interfaces.ITransaction,
|
||||
interfaces.ITransactionDeprecated)
|
||||
class Transaction(object):
|
||||
|
||||
|
||||
|
||||
# Assign an index to each savepoint so we can invalidate later savepoints
|
||||
# on rollback. The first index assigned is 1, and it goes up by 1 each
|
||||
# time.
|
||||
_savepoint_index = 0
|
||||
|
||||
# If savepoints are used, keep a weak key dict of them. This maps a
|
||||
# savepoint to its index (see above).
|
||||
_savepoint2index = None
|
||||
|
||||
# Meta data. extended_info is also metadata, but is initialized to an
|
||||
# emtpy dict in __init__.
|
||||
_user = u""
|
||||
_description = u""
|
||||
|
||||
def __init__(self, synchronizers=None, manager=None):
|
||||
self.status = Status.ACTIVE
|
||||
# List of resource managers, e.g. MultiObjectResourceAdapters.
|
||||
self._resources = []
|
||||
|
||||
# Weak set of synchronizer objects to call.
|
||||
if synchronizers is None:
|
||||
synchronizers = WeakSet()
|
||||
self._synchronizers = synchronizers
|
||||
|
||||
self._manager = manager
|
||||
|
||||
# _adapters: Connection/_p_jar -> MultiObjectResourceAdapter[Sub]
|
||||
self._adapters = {}
|
||||
self._voted = {} # id(Connection) -> boolean, True if voted
|
||||
# _voted and other dictionaries use the id() of the resource
|
||||
# manager as a key, because we can't guess whether the actual
|
||||
# resource managers will be safe to use as dict keys.
|
||||
|
||||
# The user, description, and extension attributes are accessed
|
||||
# directly by storages, leading underscore notwithstanding.
|
||||
self.extension = {}
|
||||
|
||||
self.log = _makeLogger()
|
||||
self.log.debug("new transaction")
|
||||
|
||||
# If a commit fails, the traceback is saved in _failure_traceback.
|
||||
# If another attempt is made to commit, TransactionFailedError is
|
||||
# raised, incorporating this traceback.
|
||||
self._failure_traceback = None
|
||||
|
||||
# List of (hook, args, kws) tuples added by addBeforeCommitHook().
|
||||
self._before_commit = []
|
||||
|
||||
# List of (hook, args, kws) tuples added by addAfterCommitHook().
|
||||
self._after_commit = []
|
||||
|
||||
@property
|
||||
def _extension(self):
|
||||
# for backward compatibility, since most clients used this
|
||||
# absent any formal API.
|
||||
return self.extension
|
||||
|
||||
@_extension.setter
|
||||
def _extension(self, v):
|
||||
self.extension = v
|
||||
|
||||
@property
|
||||
def user(self):
|
||||
return self._user
|
||||
|
||||
@user.setter
|
||||
def user(self, v):
|
||||
if v is None:
|
||||
raise ValueError("user must not be None")
|
||||
self._user = text_or_warn(v)
|
||||
|
||||
@property
|
||||
def description(self):
|
||||
return self._description
|
||||
|
||||
@description.setter
|
||||
def description(self, v):
|
||||
if v is not None:
|
||||
self._description = text_or_warn(v)
|
||||
|
||||
def isDoomed(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
return self.status is Status.DOOMED
|
||||
|
||||
def doom(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if self.status is not Status.DOOMED:
|
||||
if self.status is not Status.ACTIVE:
|
||||
# should not doom transactions in the middle,
|
||||
# or after, a commit
|
||||
raise ValueError('non-doomable')
|
||||
self.status = Status.DOOMED
|
||||
|
||||
# Raise TransactionFailedError, due to commit()/join()/register()
|
||||
# getting called when the current transaction has already suffered
|
||||
# a commit/savepoint failure.
|
||||
def _prior_operation_failed(self):
|
||||
assert self._failure_traceback is not None
|
||||
raise TransactionFailedError("An operation previously failed, "
|
||||
"with traceback:\n\n%s" %
|
||||
self._failure_traceback.getvalue())
|
||||
|
||||
def join(self, resource):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if self.status is Status.COMMITFAILED:
|
||||
self._prior_operation_failed() # doesn't return
|
||||
|
||||
if (self.status is not Status.ACTIVE and
|
||||
self.status is not Status.DOOMED):
|
||||
# TODO: Should it be possible to join a committing transaction?
|
||||
# I think some users want it.
|
||||
raise ValueError("expected txn status %r or %r, but it's %r" % (
|
||||
Status.ACTIVE, Status.DOOMED, self.status))
|
||||
# TODO: the prepare check is a bit of a hack, perhaps it would
|
||||
# be better to use interfaces. If this is a ZODB4-style
|
||||
# resource manager, it needs to be adapted, too.
|
||||
if myhasattr(resource, "prepare"):
|
||||
# TODO: deprecate 3.6
|
||||
resource = DataManagerAdapter(resource)
|
||||
self._resources.append(resource)
|
||||
|
||||
if self._savepoint2index:
|
||||
# A data manager has joined a transaction *after* a savepoint
|
||||
# was created. A couple of things are different in this case:
|
||||
#
|
||||
# 1. We need to add its savepoint to all previous savepoints.
|
||||
# so that if they are rolled back, we roll this one back too.
|
||||
#
|
||||
# 2. We don't actually need to ask the data manager for a
|
||||
# savepoint: because it's just joining, we can just abort it to
|
||||
# roll back to the current state, so we simply use an
|
||||
# AbortSavepoint.
|
||||
datamanager_savepoint = AbortSavepoint(resource, self)
|
||||
for transaction_savepoint in self._savepoint2index.keys():
|
||||
transaction_savepoint._savepoints.append(
|
||||
datamanager_savepoint)
|
||||
|
||||
def _unjoin(self, resource):
|
||||
# Leave a transaction because a savepoint was rolled back on a resource
|
||||
# that joined later.
|
||||
|
||||
# Don't use remove. We don't want to assume anything about __eq__.
|
||||
self._resources = [r for r in self._resources if r is not resource]
|
||||
|
||||
def savepoint(self, optimistic=False):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if self.status is Status.COMMITFAILED:
|
||||
self._prior_operation_failed() # doesn't return, it raises
|
||||
|
||||
try:
|
||||
savepoint = Savepoint(self, optimistic, *self._resources)
|
||||
except:
|
||||
self._cleanup(self._resources)
|
||||
self._saveAndRaiseCommitishError() # reraises!
|
||||
|
||||
if self._savepoint2index is None:
|
||||
self._savepoint2index = weakref.WeakKeyDictionary()
|
||||
self._savepoint_index += 1
|
||||
self._savepoint2index[savepoint] = self._savepoint_index
|
||||
|
||||
return savepoint
|
||||
|
||||
# Remove and invalidate all savepoints we know about with an index
|
||||
# larger than `savepoint`'s. This is what's needed when a rollback
|
||||
# _to_ `savepoint` is done.
|
||||
def _remove_and_invalidate_after(self, savepoint):
|
||||
savepoint2index = self._savepoint2index
|
||||
index = savepoint2index[savepoint]
|
||||
# use list(items()) to make copy to avoid mutating while iterating
|
||||
for savepoint, i in list(savepoint2index.items()):
|
||||
if i > index:
|
||||
savepoint.transaction = None # invalidate
|
||||
del savepoint2index[savepoint]
|
||||
|
||||
# Invalidate and forget about all savepoints.
|
||||
def _invalidate_all_savepoints(self):
|
||||
for savepoint in self._savepoint2index.keys():
|
||||
savepoint.transaction = None # invalidate
|
||||
self._savepoint2index.clear()
|
||||
|
||||
|
||||
def register(self, obj):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
# The old way of registering transaction participants.
|
||||
#
|
||||
# register() is passed either a persisent object or a
|
||||
# resource manager like the ones defined in ZODB.DB.
|
||||
# If it is passed a persistent object, that object should
|
||||
# be stored when the transaction commits. For other
|
||||
# objects, the object implements the standard two-phase
|
||||
# commit protocol.
|
||||
manager = getattr(obj, "_p_jar", obj)
|
||||
if manager is None:
|
||||
raise ValueError("Register with no manager")
|
||||
adapter = self._adapters.get(manager)
|
||||
if adapter is None:
|
||||
adapter = MultiObjectResourceAdapter(manager)
|
||||
adapter.objects.append(obj)
|
||||
self._adapters[manager] = adapter
|
||||
self.join(adapter)
|
||||
else:
|
||||
# TODO: comment out this expensive assert later
|
||||
# Use id() to guard against proxies.
|
||||
assert id(obj) not in map(id, adapter.objects)
|
||||
adapter.objects.append(obj)
|
||||
|
||||
def commit(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if self.status is Status.DOOMED:
|
||||
raise interfaces.DoomedTransaction(
|
||||
'transaction doomed, cannot commit')
|
||||
|
||||
if self._savepoint2index:
|
||||
self._invalidate_all_savepoints()
|
||||
|
||||
if self.status is Status.COMMITFAILED:
|
||||
self._prior_operation_failed() # doesn't return
|
||||
|
||||
self._callBeforeCommitHooks()
|
||||
|
||||
self._synchronizers.map(lambda s: s.beforeCompletion(self))
|
||||
self.status = Status.COMMITTING
|
||||
|
||||
try:
|
||||
self._commitResources()
|
||||
self.status = Status.COMMITTED
|
||||
except:
|
||||
t = None
|
||||
v = None
|
||||
tb = None
|
||||
try:
|
||||
t, v, tb = self._saveAndGetCommitishError()
|
||||
self._callAfterCommitHooks(status=False)
|
||||
reraise(t, v, tb)
|
||||
finally:
|
||||
del t, v, tb
|
||||
else:
|
||||
self._free()
|
||||
self._synchronizers.map(lambda s: s.afterCompletion(self))
|
||||
self._callAfterCommitHooks(status=True)
|
||||
self.log.debug("commit")
|
||||
|
||||
def _saveAndGetCommitishError(self):
|
||||
self.status = Status.COMMITFAILED
|
||||
# Save the traceback for TransactionFailedError.
|
||||
ft = self._failure_traceback = _makeTracebackBuffer()
|
||||
t = None
|
||||
v = None
|
||||
tb = None
|
||||
try:
|
||||
t, v, tb = sys.exc_info()
|
||||
# Record how we got into commit().
|
||||
traceback.print_stack(sys._getframe(1), None, ft)
|
||||
# Append the stack entries from here down to the exception.
|
||||
traceback.print_tb(tb, None, ft)
|
||||
# Append the exception type and value.
|
||||
ft.writelines(traceback.format_exception_only(t, v))
|
||||
return t, v, tb
|
||||
finally:
|
||||
del t, v, tb
|
||||
|
||||
def _saveAndRaiseCommitishError(self):
|
||||
t = None
|
||||
v = None
|
||||
tb = None
|
||||
try:
|
||||
t, v, tb = self._saveAndGetCommitishError()
|
||||
reraise(t, v, tb)
|
||||
finally:
|
||||
del t, v, tb
|
||||
|
||||
def getBeforeCommitHooks(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
return iter(self._before_commit)
|
||||
|
||||
def addBeforeCommitHook(self, hook, args=(), kws=None):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if kws is None:
|
||||
kws = {}
|
||||
self._before_commit.append((hook, tuple(args), kws))
|
||||
|
||||
def _callBeforeCommitHooks(self):
|
||||
# Call all hooks registered, allowing further registrations
|
||||
# during processing. Note that calls to addBeforeCommitHook() may
|
||||
# add additional hooks while hooks are running, and iterating over a
|
||||
# growing list is well-defined in Python.
|
||||
for hook, args, kws in self._before_commit:
|
||||
hook(*args, **kws)
|
||||
self._before_commit = []
|
||||
|
||||
def getAfterCommitHooks(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
return iter(self._after_commit)
|
||||
|
||||
def addAfterCommitHook(self, hook, args=(), kws=None):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if kws is None:
|
||||
kws = {}
|
||||
self._after_commit.append((hook, tuple(args), kws))
|
||||
|
||||
def _callAfterCommitHooks(self, status=True):
|
||||
# Avoid to abort anything at the end if no hooks are registred.
|
||||
if not self._after_commit:
|
||||
return
|
||||
# Call all hooks registered, allowing further registrations
|
||||
# during processing. Note that calls to addAterCommitHook() may
|
||||
# add additional hooks while hooks are running, and iterating over a
|
||||
# growing list is well-defined in Python.
|
||||
for hook, args, kws in self._after_commit:
|
||||
# The first argument passed to the hook is a Boolean value,
|
||||
# true if the commit succeeded, or false if the commit aborted.
|
||||
try:
|
||||
hook(status, *args, **kws)
|
||||
except:
|
||||
# We need to catch the exceptions if we want all hooks
|
||||
# to be called
|
||||
self.log.error("Error in after commit hook exec in %s ",
|
||||
hook, exc_info=sys.exc_info())
|
||||
# The transaction is already committed. It must not have
|
||||
# further effects after the commit.
|
||||
for rm in self._resources:
|
||||
try:
|
||||
rm.abort(self)
|
||||
except:
|
||||
# XXX should we take further actions here ?
|
||||
self.log.error("Error in abort() on manager %s",
|
||||
rm, exc_info=sys.exc_info())
|
||||
self._after_commit = []
|
||||
self._before_commit = []
|
||||
|
||||
def _commitResources(self):
|
||||
# Execute the two-phase commit protocol.
|
||||
|
||||
L = list(self._resources)
|
||||
L.sort(key=rm_key)
|
||||
try:
|
||||
for rm in L:
|
||||
rm.tpc_begin(self)
|
||||
for rm in L:
|
||||
rm.commit(self)
|
||||
self.log.debug("commit %r", rm)
|
||||
for rm in L:
|
||||
rm.tpc_vote(self)
|
||||
self._voted[id(rm)] = True
|
||||
|
||||
try:
|
||||
for rm in L:
|
||||
rm.tpc_finish(self)
|
||||
except:
|
||||
# TODO: do we need to make this warning stronger?
|
||||
# TODO: It would be nice if the system could be configured
|
||||
# to stop committing transactions at this point.
|
||||
self.log.critical("A storage error occurred during the second "
|
||||
"phase of the two-phase commit. Resources "
|
||||
"may be in an inconsistent state.")
|
||||
raise
|
||||
except:
|
||||
# If an error occurs committing a transaction, we try
|
||||
# to revert the changes in each of the resource managers.
|
||||
t, v, tb = sys.exc_info()
|
||||
try:
|
||||
try:
|
||||
self._cleanup(L)
|
||||
finally:
|
||||
self._synchronizers.map(lambda s: s.afterCompletion(self))
|
||||
reraise(t, v, tb)
|
||||
finally:
|
||||
del t, v, tb
|
||||
|
||||
def _cleanup(self, L):
|
||||
# Called when an exception occurs during tpc_vote or tpc_finish.
|
||||
for rm in L:
|
||||
if id(rm) not in self._voted:
|
||||
try:
|
||||
rm.abort(self)
|
||||
except Exception:
|
||||
self.log.error("Error in abort() on manager %s",
|
||||
rm, exc_info=sys.exc_info())
|
||||
for rm in L:
|
||||
try:
|
||||
rm.tpc_abort(self)
|
||||
except Exception:
|
||||
self.log.error("Error in tpc_abort() on manager %s",
|
||||
rm, exc_info=sys.exc_info())
|
||||
|
||||
def _free(self):
|
||||
# Called when the transaction has been committed or aborted
|
||||
# to break references---this transaction object will not be returned
|
||||
# as the current transaction from its manager after this, and all
|
||||
# IDatamanager objects joined to it will forgotten
|
||||
if self._manager:
|
||||
self._manager.free(self)
|
||||
|
||||
if hasattr(self, '_data'):
|
||||
delattr(self, '_data')
|
||||
|
||||
del self._resources[:]
|
||||
|
||||
def data(self, ob):
|
||||
try:
|
||||
data = self._data
|
||||
except AttributeError:
|
||||
raise KeyError(ob)
|
||||
|
||||
try:
|
||||
return data[id(ob)]
|
||||
except KeyError:
|
||||
raise KeyError(ob)
|
||||
|
||||
def set_data(self, ob, ob_data):
|
||||
try:
|
||||
data = self._data
|
||||
except AttributeError:
|
||||
data = self._data = {}
|
||||
|
||||
data[id(ob)] = ob_data
|
||||
|
||||
def abort(self):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if self._savepoint2index:
|
||||
self._invalidate_all_savepoints()
|
||||
|
||||
self._synchronizers.map(lambda s: s.beforeCompletion(self))
|
||||
|
||||
try:
|
||||
|
||||
t = None
|
||||
v = None
|
||||
tb = None
|
||||
|
||||
for rm in self._resources:
|
||||
try:
|
||||
rm.abort(self)
|
||||
except:
|
||||
if tb is None:
|
||||
t, v, tb = sys.exc_info()
|
||||
self.log.error("Failed to abort resource manager: %s",
|
||||
rm, exc_info=sys.exc_info())
|
||||
|
||||
self._free()
|
||||
|
||||
self._synchronizers.map(lambda s: s.afterCompletion(self))
|
||||
|
||||
self.log.debug("abort")
|
||||
|
||||
if tb is not None:
|
||||
reraise(t, v, tb)
|
||||
finally:
|
||||
del t, v, tb
|
||||
|
||||
def note(self, text):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
if text is not None:
|
||||
text = text_or_warn(text).strip()
|
||||
if self.description:
|
||||
self.description += u"\n" + text
|
||||
else:
|
||||
self.description = text
|
||||
|
||||
def setUser(self, user_name, path=u"/"):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
self.user = u"%s %s" % (text_or_warn(path), text_or_warn(user_name))
|
||||
|
||||
def setExtendedInfo(self, name, value):
|
||||
""" See ITransaction.
|
||||
"""
|
||||
self.extension[name] = value
|
||||
|
||||
def isRetryableError(self, error):
|
||||
return self._manager._retryable(type(error), error)
|
||||
|
||||
|
||||
# TODO: We need a better name for the adapters.
|
||||
|
||||
|
||||
class MultiObjectResourceAdapter(object):
|
||||
"""Adapt the old-style register() call to the new-style join().
|
||||
|
||||
With join(), a resource manager like a Connection registers with
|
||||
the transaction manager. With register(), an individual object
|
||||
is passed to register().
|
||||
"""
|
||||
def __init__(self, jar):
|
||||
self.manager = jar
|
||||
self.objects = []
|
||||
self.ncommitted = 0
|
||||
|
||||
def __repr__(self):
|
||||
return "<%s for %s at %s>" % (self.__class__.__name__,
|
||||
self.manager, id(self))
|
||||
|
||||
def sortKey(self):
|
||||
return self.manager.sortKey()
|
||||
|
||||
def tpc_begin(self, txn):
|
||||
self.manager.tpc_begin(txn)
|
||||
|
||||
def tpc_finish(self, txn):
|
||||
self.manager.tpc_finish(txn)
|
||||
|
||||
def tpc_abort(self, txn):
|
||||
self.manager.tpc_abort(txn)
|
||||
|
||||
def commit(self, txn):
|
||||
for o in self.objects:
|
||||
self.manager.commit(o, txn)
|
||||
self.ncommitted += 1
|
||||
|
||||
def tpc_vote(self, txn):
|
||||
self.manager.tpc_vote(txn)
|
||||
|
||||
def abort(self, txn):
|
||||
t = None
|
||||
v = None
|
||||
tb = None
|
||||
try:
|
||||
for o in self.objects:
|
||||
try:
|
||||
self.manager.abort(o, txn)
|
||||
except:
|
||||
# Capture the first exception and re-raise it after
|
||||
# aborting all the other objects.
|
||||
if tb is None:
|
||||
t, v, tb = sys.exc_info()
|
||||
txn.log.error("Failed to abort object: %s",
|
||||
object_hint(o), exc_info=sys.exc_info())
|
||||
|
||||
if tb is not None:
|
||||
reraise(t, v, tb)
|
||||
finally:
|
||||
del t, v, tb
|
||||
|
||||
|
||||
def rm_key(rm):
|
||||
func = getattr(rm, 'sortKey', None)
|
||||
if func is not None:
|
||||
return func()
|
||||
|
||||
def object_hint(o):
|
||||
"""Return a string describing the object.
|
||||
|
||||
This function does not raise an exception.
|
||||
"""
|
||||
# We should always be able to get __class__.
|
||||
klass = o.__class__.__name__
|
||||
# oid would be great, but maybe this isn't a persistent object.
|
||||
oid = getattr(o, "_p_oid", _marker)
|
||||
if oid is not _marker:
|
||||
oid = oid_repr(oid)
|
||||
else:
|
||||
oid = 'None'
|
||||
return "%s oid=%s" % (klass, oid)
|
||||
|
||||
def oid_repr(oid):
|
||||
if isinstance(oid, str) and len(oid) == 8:
|
||||
# Convert to hex and strip leading zeroes.
|
||||
as_hex = native_(
|
||||
binascii.hexlify(bytes_(oid, 'ascii')), 'ascii').lstrip('0')
|
||||
# Ensure two characters per input byte.
|
||||
if len(as_hex) & 1:
|
||||
as_hex = '0' + as_hex
|
||||
elif as_hex == '':
|
||||
as_hex = '00'
|
||||
return '0x' + as_hex
|
||||
else:
|
||||
return repr(oid)
|
||||
|
||||
|
||||
# TODO: deprecate for 3.6.
|
||||
class DataManagerAdapter(object):
|
||||
"""Adapt zodb 4-style data managers to zodb3 style
|
||||
|
||||
Adapt transaction.interfaces.IDataManager to
|
||||
ZODB.interfaces.IPureDatamanager
|
||||
"""
|
||||
|
||||
# Note that it is pretty important that this does not have a _p_jar
|
||||
# attribute. This object will be registered with a zodb3 TM, which
|
||||
# will then try to get a _p_jar from it, using it as the default.
|
||||
# (Objects without a _p_jar are their own data managers.)
|
||||
|
||||
def __init__(self, datamanager):
|
||||
self._datamanager = datamanager
|
||||
|
||||
# TODO: I'm not sure why commit() doesn't do anything
|
||||
|
||||
def commit(self, transaction):
|
||||
# We don't do anything here because ZODB4-style data managers
|
||||
# didn't have a separate commit step
|
||||
pass
|
||||
|
||||
def abort(self, transaction):
|
||||
self._datamanager.abort(transaction)
|
||||
|
||||
def tpc_begin(self, transaction):
|
||||
# We don't do anything here because ZODB4-style data managers
|
||||
# didn't have a separate tpc_begin step
|
||||
pass
|
||||
|
||||
def tpc_abort(self, transaction):
|
||||
self._datamanager.abort(transaction)
|
||||
|
||||
def tpc_finish(self, transaction):
|
||||
self._datamanager.commit(transaction)
|
||||
|
||||
def tpc_vote(self, transaction):
|
||||
self._datamanager.prepare(transaction)
|
||||
|
||||
def sortKey(self):
|
||||
return self._datamanager.sortKey()
|
||||
|
||||
|
||||
@implementer(interfaces.ISavepoint)
|
||||
class Savepoint:
|
||||
"""Transaction savepoint.
|
||||
|
||||
Transaction savepoints coordinate savepoints for data managers
|
||||
participating in a transaction.
|
||||
"""
|
||||
|
||||
def __init__(self, transaction, optimistic, *resources):
|
||||
self.transaction = transaction
|
||||
self._savepoints = savepoints = []
|
||||
|
||||
for datamanager in resources:
|
||||
try:
|
||||
savepoint = datamanager.savepoint
|
||||
except AttributeError:
|
||||
if not optimistic:
|
||||
raise TypeError("Savepoints unsupported", datamanager)
|
||||
savepoint = NoRollbackSavepoint(datamanager)
|
||||
else:
|
||||
savepoint = savepoint()
|
||||
|
||||
savepoints.append(savepoint)
|
||||
|
||||
@property
|
||||
def valid(self):
|
||||
return self.transaction is not None
|
||||
|
||||
def rollback(self):
|
||||
""" See ISavepoint.
|
||||
"""
|
||||
transaction = self.transaction
|
||||
if transaction is None:
|
||||
raise interfaces.InvalidSavepointRollbackError(
|
||||
'invalidated by a later savepoint')
|
||||
transaction._remove_and_invalidate_after(self)
|
||||
|
||||
try:
|
||||
for savepoint in self._savepoints:
|
||||
savepoint.rollback()
|
||||
except:
|
||||
# Mark the transaction as failed.
|
||||
transaction._saveAndRaiseCommitishError() # reraises!
|
||||
|
||||
|
||||
class AbortSavepoint:
|
||||
|
||||
def __init__(self, datamanager, transaction):
|
||||
self.datamanager = datamanager
|
||||
self.transaction = transaction
|
||||
|
||||
def rollback(self):
|
||||
self.datamanager.abort(self.transaction)
|
||||
self.transaction._unjoin(self.datamanager)
|
||||
|
||||
|
||||
class NoRollbackSavepoint:
|
||||
|
||||
def __init__(self, datamanager):
|
||||
self.datamanager = datamanager
|
||||
|
||||
def rollback(self):
|
||||
raise TypeError("Savepoints unsupported", self.datamanager)
|
||||
|
||||
def text_or_warn(s):
|
||||
if isinstance(s, text_type):
|
||||
return s
|
||||
|
||||
warnings.warn("Expected text", DeprecationWarning, stacklevel=3)
|
||||
if isinstance(s, bytes):
|
||||
return s.decode('utf-8', 'replace')
|
||||
else:
|
||||
return text_type(s)
|
||||
@@ -0,0 +1,586 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2001, 2002 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
##############################################################################
|
||||
|
||||
from zope.interface import Attribute
|
||||
from zope.interface import Interface
|
||||
|
||||
class ITransactionManager(Interface):
|
||||
"""An object that manages a sequence of transactions.
|
||||
|
||||
Applications use transaction managers to establish transaction boundaries.
|
||||
"""
|
||||
|
||||
explicit = Attribute(
|
||||
"""Explicit mode indicator.
|
||||
|
||||
This is true if the transaction manager is in explicit mode.
|
||||
In explicit mode, transactions must be begun explicitly, by
|
||||
calling ``begin()`` and ended explicitly by calling
|
||||
``commit()`` or ``abort()``.
|
||||
""")
|
||||
|
||||
|
||||
def begin():
|
||||
"""Explicitly begin and return a new transaction.
|
||||
|
||||
If an existing transaction is in progress and the transaction
|
||||
manager not in explicit mode, the previous transaction will be
|
||||
aborted. If an existing transaction is in progress and the
|
||||
transaction manager is in explicit mode, an
|
||||
``AlreadyInTransaction`` exception will be raised..
|
||||
|
||||
The ``newTransaction`` method of registered synchronizers is called,
|
||||
passing the new transaction object.
|
||||
|
||||
Note that when not in explicit mode, transactions may be
|
||||
started implicitly without calling ``begin``. In that case,
|
||||
``newTransaction`` isn't called because the transaction
|
||||
manager doesn't know when to call it. The transaction is
|
||||
likely to have begun long before the transaction manager is
|
||||
involved. (Conceivably the ``commit`` and ``abort`` methods
|
||||
could call ``begin``, but they don't.)
|
||||
"""
|
||||
|
||||
def get():
|
||||
"""Get the current transaction.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def commit():
|
||||
"""Commit the current transaction.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def abort():
|
||||
"""Abort the current transaction.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def doom():
|
||||
"""Doom the current transaction.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def isDoomed():
|
||||
"""Returns True if the current transaction is doomed, otherwise False.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def savepoint(optimistic=False):
|
||||
"""Create a savepoint from the current transaction.
|
||||
|
||||
If the optimistic argument is true, then data managers that
|
||||
don't support savepoints can be used, but an error will be
|
||||
raised if the savepoint is rolled back.
|
||||
|
||||
An ISavepoint object is returned.
|
||||
|
||||
In explicit mode, if a transaction hasn't begun, a
|
||||
``NoTransaction`` exception will be raised.
|
||||
"""
|
||||
|
||||
def registerSynch(synch):
|
||||
"""Register an ISynchronizer.
|
||||
|
||||
Synchronizers are notified about some major events in a transaction's
|
||||
life. See ISynchronizer for details.
|
||||
|
||||
If a synchronizer registers while there is an active
|
||||
transaction, its newTransaction method will be called with the
|
||||
active transaction.
|
||||
"""
|
||||
|
||||
def unregisterSynch(synch):
|
||||
"""Unregister an ISynchronizer.
|
||||
|
||||
Synchronizers are notified about some major events in a transaction's
|
||||
life. See ISynchronizer for details.
|
||||
"""
|
||||
|
||||
def clearSynchs():
|
||||
"""Unregister all registered ISynchronizers.
|
||||
|
||||
This exists to support test cleanup/initialization
|
||||
"""
|
||||
|
||||
def registeredSynchs():
|
||||
"""Determine if any ISynchronizers are registered.
|
||||
|
||||
Return true if any are registered, and return False otherwise.
|
||||
|
||||
This exists to support test cleanup/initialization
|
||||
"""
|
||||
|
||||
class ITransaction(Interface):
|
||||
"""Object representing a running transaction.
|
||||
|
||||
Objects with this interface may represent different transactions
|
||||
during their lifetime (.begin() can be called to start a new
|
||||
transaction using the same instance, although that example is
|
||||
deprecated and will go away in ZODB 3.6).
|
||||
"""
|
||||
|
||||
user = Attribute(
|
||||
"""A user name associated with the transaction.
|
||||
|
||||
The format of the user name is defined by the application. The value
|
||||
is text (unicode). Storages record the user value, as meta-data,
|
||||
when a transaction commits.
|
||||
|
||||
A storage may impose a limit on the size of the value; behavior is
|
||||
undefined if such a limit is exceeded (for example, a storage may
|
||||
raise an exception, or truncate the value).
|
||||
""")
|
||||
|
||||
description = Attribute(
|
||||
"""A textual description of the transaction.
|
||||
|
||||
The value is text (unicode). Method note() is the intended
|
||||
way to set the value. Storages record the description, as meta-data,
|
||||
when a transaction commits.
|
||||
|
||||
A storage may impose a limit on the size of the description; behavior
|
||||
is undefined if such a limit is exceeded (for example, a storage may
|
||||
raise an exception, or truncate the value).
|
||||
""")
|
||||
|
||||
extension = Attribute(
|
||||
"A dictionary containing application-defined metadata.")
|
||||
|
||||
def commit():
|
||||
"""Finalize the transaction.
|
||||
|
||||
This executes the two-phase commit algorithm for all
|
||||
IDataManager objects associated with the transaction.
|
||||
"""
|
||||
|
||||
def abort():
|
||||
"""Abort the transaction.
|
||||
|
||||
This is called from the application. This can only be called
|
||||
before the two-phase commit protocol has been started.
|
||||
"""
|
||||
|
||||
def doom():
|
||||
"""Doom the transaction.
|
||||
|
||||
Dooms the current transaction. This will cause
|
||||
DoomedTransactionException to be raised on any attempt to commit the
|
||||
transaction.
|
||||
|
||||
Otherwise the transaction will behave as if it was active.
|
||||
"""
|
||||
|
||||
def savepoint(optimistic=False):
|
||||
"""Create a savepoint.
|
||||
|
||||
If the optimistic argument is true, then data managers that don't
|
||||
support savepoints can be used, but an error will be raised if the
|
||||
savepoint is rolled back.
|
||||
|
||||
An ISavepoint object is returned.
|
||||
"""
|
||||
|
||||
def join(datamanager):
|
||||
"""Add a data manager to the transaction.
|
||||
|
||||
`datamanager` must provide the transactions.interfaces.IDataManager
|
||||
interface.
|
||||
"""
|
||||
|
||||
def note(text):
|
||||
"""Add text (unicode) to the transaction description.
|
||||
|
||||
This modifies the `.description` attribute; see its docs for more
|
||||
detail. First surrounding whitespace is stripped from `text`. If
|
||||
`.description` is currently an empty string, then the stripped text
|
||||
becomes its value, else two newlines and the stripped text are
|
||||
appended to `.description`.
|
||||
"""
|
||||
|
||||
def setExtendedInfo(name, value):
|
||||
"""Add extension data to the transaction.
|
||||
|
||||
name
|
||||
is the text (unicode) name of the extension property to set
|
||||
|
||||
value
|
||||
must be picklable and json serializable (not an instance).
|
||||
|
||||
Multiple calls may be made to set multiple extension
|
||||
properties, provided the names are distinct.
|
||||
|
||||
Storages record the extension data, as meta-data, when a transaction
|
||||
commits.
|
||||
|
||||
A storage may impose a limit on the size of extension data; behavior
|
||||
is undefined if such a limit is exceeded (for example, a storage may
|
||||
raise an exception, or remove `<name, value>` pairs).
|
||||
"""
|
||||
|
||||
def addBeforeCommitHook(hook, args=(), kws=None):
|
||||
"""Register a hook to call before the transaction is committed.
|
||||
|
||||
The specified hook function will be called after the transaction's
|
||||
commit method has been called, but before the commit process has been
|
||||
started. The hook will be passed the specified positional (`args`)
|
||||
and keyword (`kws`) arguments. `args` is a sequence of positional
|
||||
arguments to be passed, defaulting to an empty tuple (no positional
|
||||
arguments are passed). `kws` is a dictionary of keyword argument
|
||||
names and values to be passed, or the default None (no keyword
|
||||
arguments are passed).
|
||||
|
||||
Multiple hooks can be registered and will be called in the order they
|
||||
were registered (first registered, first called). This method can
|
||||
also be called from a hook: an executing hook can register more
|
||||
hooks. Applications should take care to avoid creating infinite loops
|
||||
by recursively registering hooks.
|
||||
|
||||
Hooks are called only for a top-level commit. A
|
||||
savepoint creation does not call any hooks. If the
|
||||
transaction is aborted, hooks are not called, and are discarded.
|
||||
Calling a hook "consumes" its registration too: hook registrations
|
||||
do not persist across transactions. If it's desired to call the same
|
||||
hook on every transaction commit, then addBeforeCommitHook() must be
|
||||
called with that hook during every transaction; in such a case
|
||||
consider registering a synchronizer object via a TransactionManager's
|
||||
registerSynch() method instead.
|
||||
"""
|
||||
|
||||
def getBeforeCommitHooks():
|
||||
"""Return iterable producing the registered addBeforeCommit hooks.
|
||||
|
||||
A triple (hook, args, kws) is produced for each registered hook.
|
||||
The hooks are produced in the order in which they would be invoked
|
||||
by a top-level transaction commit.
|
||||
"""
|
||||
|
||||
def addAfterCommitHook(hook, args=(), kws=None):
|
||||
"""Register a hook to call after a transaction commit attempt.
|
||||
|
||||
The specified hook function will be called after the transaction
|
||||
commit succeeds or aborts. The first argument passed to the hook
|
||||
is a Boolean value, true if the commit succeeded, or false if the
|
||||
commit aborted. `args` specifies additional positional, and `kws`
|
||||
keyword, arguments to pass to the hook. `args` is a sequence of
|
||||
positional arguments to be passed, defaulting to an empty tuple
|
||||
(only the true/false success argument is passed). `kws` is a
|
||||
dictionary of keyword argument names and values to be passed, or
|
||||
the default None (no keyword arguments are passed).
|
||||
|
||||
Multiple hooks can be registered and will be called in the order they
|
||||
were registered (first registered, first called). This method can
|
||||
also be called from a hook: an executing hook can register more
|
||||
hooks. Applications should take care to avoid creating infinite loops
|
||||
by recursively registering hooks.
|
||||
|
||||
Hooks are called only for a top-level commit. A
|
||||
savepoint creation does not call any hooks. Calling a
|
||||
hook "consumes" its registration: hook registrations do not
|
||||
persist across transactions. If it's desired to call the same
|
||||
hook on every transaction commit, then addAfterCommitHook() must be
|
||||
called with that hook during every transaction; in such a case
|
||||
consider registering a synchronizer object via a TransactionManager's
|
||||
registerSynch() method instead.
|
||||
"""
|
||||
|
||||
def getAfterCommitHooks():
|
||||
"""Return iterable producing the registered addAfterCommit hooks.
|
||||
|
||||
A triple (hook, args, kws) is produced for each registered hook.
|
||||
The hooks are produced in the order in which they would be invoked
|
||||
by a top-level transaction commit.
|
||||
"""
|
||||
|
||||
def set_data(ob, data):
|
||||
"""Hold data on behalf of an object
|
||||
|
||||
For objects such as data managers or their subobjects that
|
||||
work with multiple transactions, it's convenient to store
|
||||
transaction-specific data on the transaction itself. The
|
||||
transaction knows nothing about the data, but simply holds it
|
||||
on behalf of the object.
|
||||
|
||||
The object passed should be the object that needs the data, as
|
||||
opposed to simple object like a string. (Internally, the id of
|
||||
the object is used as the key.)
|
||||
"""
|
||||
|
||||
def data(ob):
|
||||
"""Retrieve data held on behalf of an object.
|
||||
|
||||
See set_data.
|
||||
"""
|
||||
|
||||
def isRetryableError(error):
|
||||
"""Determine if the error is retryable.
|
||||
|
||||
Return true if any joined IRetryDataManager considers the error
|
||||
transient. Such errors may occur due to concurrency issues in the
|
||||
underlying storage engine.
|
||||
|
||||
"""
|
||||
|
||||
class ITransactionDeprecated(Interface):
|
||||
"""Deprecated parts of the transaction API."""
|
||||
|
||||
def begin(info=None):
|
||||
"""Begin a new transaction.
|
||||
|
||||
If the transaction is in progress, it is aborted and a new
|
||||
transaction is started using the same transaction object.
|
||||
"""
|
||||
|
||||
# TODO: deprecate this for 3.6.
|
||||
def register(object):
|
||||
"""Register the given object for transaction control."""
|
||||
|
||||
|
||||
class IDataManager(Interface):
|
||||
"""Objects that manage transactional storage.
|
||||
|
||||
These objects may manage data for other objects, or they may manage
|
||||
non-object storages, such as relational databases. For example,
|
||||
a ZODB.Connection.
|
||||
|
||||
Note that when some data is modified, that data's data manager should
|
||||
join a transaction so that data can be committed when the user commits
|
||||
the transaction.
|
||||
"""
|
||||
|
||||
transaction_manager = Attribute(
|
||||
"""The transaction manager (TM) used by this data manager.
|
||||
|
||||
This is a public attribute, intended for read-only use. The value
|
||||
is an instance of ITransactionManager, typically set by the data
|
||||
manager's constructor.
|
||||
""")
|
||||
|
||||
def abort(transaction):
|
||||
"""Abort a transaction and forget all changes.
|
||||
|
||||
Abort must be called outside of a two-phase commit.
|
||||
|
||||
Abort is called by the transaction manager to abort
|
||||
transactions that are not yet in a two-phase commit. It may
|
||||
also be called when rolling back a savepoint made before the
|
||||
data manager joined the transaction.
|
||||
|
||||
In any case, after abort is called, the data manager is no
|
||||
longer participating in the transaction. If there are new
|
||||
changes, the data manager must rejoin the transaction.
|
||||
"""
|
||||
|
||||
# Two-phase commit protocol. These methods are called by the ITransaction
|
||||
# object associated with the transaction being committed. The sequence
|
||||
# of calls normally follows this regular expression:
|
||||
# tpc_begin commit tpc_vote (tpc_finish | tpc_abort)
|
||||
|
||||
def tpc_begin(transaction):
|
||||
"""Begin commit of a transaction, starting the two-phase commit.
|
||||
|
||||
transaction is the ITransaction instance associated with the
|
||||
transaction being committed.
|
||||
"""
|
||||
|
||||
def commit(transaction):
|
||||
"""Commit modifications to registered objects.
|
||||
|
||||
Save changes to be made persistent if the transaction commits (if
|
||||
tpc_finish is called later). If tpc_abort is called later, changes
|
||||
must not persist.
|
||||
|
||||
This includes conflict detection and handling. If no conflicts or
|
||||
errors occur, the data manager should be prepared to make the
|
||||
changes persist when tpc_finish is called.
|
||||
"""
|
||||
|
||||
def tpc_vote(transaction):
|
||||
"""Verify that a data manager can commit the transaction.
|
||||
|
||||
This is the last chance for a data manager to vote 'no'. A
|
||||
data manager votes 'no' by raising an exception.
|
||||
|
||||
transaction is the ITransaction instance associated with the
|
||||
transaction being committed.
|
||||
"""
|
||||
|
||||
def tpc_finish(transaction):
|
||||
"""Indicate confirmation that the transaction is done.
|
||||
|
||||
Make all changes to objects modified by this transaction persist.
|
||||
|
||||
transaction is the ITransaction instance associated with the
|
||||
transaction being committed.
|
||||
|
||||
This should never fail. If this raises an exception, the
|
||||
database is not expected to maintain consistency; it's a
|
||||
serious error.
|
||||
"""
|
||||
|
||||
def tpc_abort(transaction):
|
||||
"""Abort a transaction.
|
||||
|
||||
This is called by a transaction manager to end a two-phase commit on
|
||||
the data manager. Abandon all changes to objects modified by this
|
||||
transaction.
|
||||
|
||||
transaction is the ITransaction instance associated with the
|
||||
transaction being committed.
|
||||
|
||||
This should never fail.
|
||||
"""
|
||||
|
||||
def sortKey():
|
||||
"""Return a key to use for ordering registered DataManagers.
|
||||
|
||||
In order to guarantee a total ordering, keys must be strings.
|
||||
|
||||
ZODB uses a global sort order to prevent deadlock when it commits
|
||||
transactions involving multiple resource managers. The resource
|
||||
manager must define a sortKey() method that provides a global ordering
|
||||
for resource managers.
|
||||
"""
|
||||
# Alternate version:
|
||||
#"""Return a consistent sort key for this connection.
|
||||
#
|
||||
#This allows ordering multiple connections that use the same storage in
|
||||
#a consistent manner. This is unique for the lifetime of a connection,
|
||||
#which is good enough to avoid ZEO deadlocks.
|
||||
#"""
|
||||
|
||||
class ISavepointDataManager(IDataManager):
|
||||
|
||||
def savepoint():
|
||||
"""Return a data-manager savepoint (IDataManagerSavepoint).
|
||||
"""
|
||||
|
||||
class IRetryDataManager(IDataManager):
|
||||
|
||||
def should_retry(exception):
|
||||
"""Return whether a given exception instance should be retried.
|
||||
|
||||
A data manager can provide this method to indicate that a a
|
||||
transaction that raised the given error should be retried.
|
||||
This method may be called by an ITransactionManager when
|
||||
considering whether to retry a failed transaction.
|
||||
"""
|
||||
|
||||
class IDataManagerSavepoint(Interface):
|
||||
"""Savepoint for data-manager changes for use in transaction savepoints.
|
||||
|
||||
Datamanager savepoints are used by, and only by, transaction savepoints.
|
||||
|
||||
Note that data manager savepoints don't have any notion of, or
|
||||
responsibility for, validity. It isn't the responsibility of
|
||||
data-manager savepoints to prevent multiple rollbacks or rollbacks after
|
||||
transaction termination. Preventing invalid savepoint rollback is the
|
||||
responsibility of transaction rollbacks. Application code should never
|
||||
use data-manager savepoints.
|
||||
"""
|
||||
|
||||
def rollback():
|
||||
"""Rollback any work done since the savepoint.
|
||||
"""
|
||||
|
||||
class ISavepoint(Interface):
|
||||
"""A transaction savepoint.
|
||||
"""
|
||||
|
||||
def rollback():
|
||||
"""Rollback any work done since the savepoint.
|
||||
|
||||
InvalidSavepointRollbackError is raised if the savepoint isn't valid.
|
||||
"""
|
||||
|
||||
valid = Attribute(
|
||||
"Boolean indicating whether the savepoint is valid")
|
||||
|
||||
class InvalidSavepointRollbackError(Exception):
|
||||
"""Attempt to rollback an invalid savepoint.
|
||||
|
||||
A savepoint may be invalid because:
|
||||
|
||||
- The surrounding transaction has committed or aborted.
|
||||
|
||||
- An earlier savepoint in the same transaction has been rolled back.
|
||||
"""
|
||||
|
||||
class ISynchronizer(Interface):
|
||||
"""Objects that participate in the transaction-boundary notification API.
|
||||
"""
|
||||
|
||||
def beforeCompletion(transaction):
|
||||
"""Hook that is called by the transaction at the start of a commit.
|
||||
"""
|
||||
|
||||
def afterCompletion(transaction):
|
||||
"""Hook that is called by the transaction after completing a commit.
|
||||
"""
|
||||
|
||||
def newTransaction(transaction):
|
||||
"""Hook that is called at the start of a transaction.
|
||||
|
||||
This hook is called when, and only when, a transaction manager's
|
||||
begin() method is called explictly.
|
||||
"""
|
||||
|
||||
class TransactionError(Exception):
|
||||
"""An error occurred due to normal transaction processing."""
|
||||
|
||||
class TransactionFailedError(TransactionError):
|
||||
"""Cannot perform an operation on a transaction that previously failed.
|
||||
|
||||
An attempt was made to commit a transaction, or to join a transaction,
|
||||
but this transaction previously raised an exception during an attempt
|
||||
to commit it. The transaction must be explicitly aborted, either by
|
||||
invoking abort() on the transaction, or begin() on its transaction
|
||||
manager.
|
||||
"""
|
||||
|
||||
class DoomedTransaction(TransactionError):
|
||||
"""A commit was attempted on a transaction that was doomed."""
|
||||
|
||||
class TransientError(TransactionError):
|
||||
"""An error has occured when performing a transaction.
|
||||
|
||||
It's possible that retrying the transaction will succeed.
|
||||
"""
|
||||
|
||||
class NoTransaction(TransactionError):
|
||||
"""No transaction has been defined
|
||||
|
||||
An application called an operation on a transaction manager that
|
||||
affects an exciting transaction, but no transaction was begun.
|
||||
The transaction manager was in explicit mode, so a new transaction
|
||||
was not explicitly created.
|
||||
"""
|
||||
|
||||
class AlreadyInTransaction(TransactionError):
|
||||
"""Attempt to create a new transaction without ending a preceding one
|
||||
|
||||
An application called ``begin()`` on a transaction manager in
|
||||
explicit mode, without committing or aborting the previous
|
||||
transaction.
|
||||
"""
|
||||
@@ -0,0 +1 @@
|
||||
#
|
||||
@@ -0,0 +1,64 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2012 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE
|
||||
#
|
||||
##############################################################################
|
||||
|
||||
|
||||
class DummyFile(object):
|
||||
def __init__(self):
|
||||
self._lines = []
|
||||
def write(self, text):
|
||||
self._lines.append(text)
|
||||
def writelines(self, lines):
|
||||
self._lines.extend(lines)
|
||||
|
||||
|
||||
class DummyLogger(object):
|
||||
def __init__(self):
|
||||
self._clear()
|
||||
def _clear(self):
|
||||
self._log = []
|
||||
def log(self, level, msg, *args, **kwargs):
|
||||
if args:
|
||||
self._log.append((level, msg % args))
|
||||
else:
|
||||
self._log.append((level, msg))
|
||||
def debug(self, msg, *args, **kw):
|
||||
self.log('debug', msg, *args, **kw)
|
||||
def error(self, msg, *args, **kw):
|
||||
self.log('error', msg, *args, **kw)
|
||||
def critical(self, msg, *args, **kw):
|
||||
self.log('critical', msg, *args, **kw)
|
||||
|
||||
|
||||
class Monkey(object):
|
||||
# context-manager for replacing module names in the scope of a test.
|
||||
def __init__(self, module, **kw):
|
||||
self.module = module
|
||||
self.to_restore = {key: getattr(module, key) for key in kw}
|
||||
for key, value in kw.items():
|
||||
setattr(module, key, value)
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, exc_type, exc_val, exc_tb):
|
||||
for key, value in self.to_restore.items():
|
||||
setattr(self.module, key, value)
|
||||
|
||||
def assertRaisesEx(e_type, checked, *args, **kw):
|
||||
# Only used in doctests
|
||||
try:
|
||||
checked(*args, **kw)
|
||||
except e_type as e:
|
||||
return e
|
||||
raise AssertionError("Didn't raise: %s" % e_type.__name__)
|
||||
@@ -0,0 +1,181 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
##############################################################################
|
||||
"""Sample objects for use in tests
|
||||
|
||||
"""
|
||||
|
||||
|
||||
class DataManager(object):
|
||||
"""Sample data manager
|
||||
|
||||
Used by the 'datamanager' chapter in the Sphinx docs.
|
||||
"""
|
||||
def __init__(self):
|
||||
self.state = 0
|
||||
self.sp = 0
|
||||
self.transaction = None
|
||||
self.delta = 0
|
||||
self.prepared = False
|
||||
|
||||
def inc(self, n=1):
|
||||
self.delta += n
|
||||
|
||||
def prepare(self, transaction):
|
||||
if self.prepared:
|
||||
raise TypeError('Already prepared')
|
||||
self._checkTransaction(transaction)
|
||||
self.prepared = True
|
||||
self.transaction = transaction
|
||||
self.state += self.delta
|
||||
|
||||
def _checkTransaction(self, transaction):
|
||||
if (transaction is not self.transaction
|
||||
and self.transaction is not None):
|
||||
raise TypeError("Transaction missmatch",
|
||||
transaction, self.transaction)
|
||||
|
||||
def abort(self, transaction):
|
||||
self._checkTransaction(transaction)
|
||||
if self.transaction is not None:
|
||||
self.transaction = None
|
||||
|
||||
if self.prepared:
|
||||
self.state -= self.delta
|
||||
self.prepared = False
|
||||
|
||||
self.delta = 0
|
||||
|
||||
def commit(self, transaction):
|
||||
if not self.prepared:
|
||||
raise TypeError('Not prepared to commit')
|
||||
self._checkTransaction(transaction)
|
||||
self.delta = 0
|
||||
self.transaction = None
|
||||
self.prepared = False
|
||||
|
||||
def savepoint(self, transaction):
|
||||
if self.prepared:
|
||||
raise AssertionError("Can't get savepoint during two-phase commit")
|
||||
self._checkTransaction(transaction)
|
||||
self.transaction = transaction
|
||||
self.sp += 1
|
||||
return Rollback(self)
|
||||
|
||||
|
||||
class Rollback(object):
|
||||
|
||||
def __init__(self, dm):
|
||||
self.dm = dm
|
||||
self.sp = dm.sp
|
||||
self.delta = dm.delta
|
||||
self.transaction = dm.transaction
|
||||
|
||||
def rollback(self):
|
||||
if self.transaction is not self.dm.transaction:
|
||||
raise TypeError("Attempt to rollback stale rollback")
|
||||
if self.dm.sp < self.sp:
|
||||
raise TypeError("Attempt to roll back to invalid save point",
|
||||
self.sp, self.dm.sp)
|
||||
self.dm.sp = self.sp
|
||||
self.dm.delta = self.delta
|
||||
|
||||
|
||||
class ResourceManager(object):
|
||||
""" Sample resource manager.
|
||||
|
||||
Used by the 'resourcemanager' chapter in the Sphinx docs.
|
||||
"""
|
||||
def __init__(self):
|
||||
self.state = 0
|
||||
self.sp = 0
|
||||
self.transaction = None
|
||||
self.delta = 0
|
||||
self.txn_state = None
|
||||
|
||||
def _check_state(self, *ok_states):
|
||||
if self.txn_state not in ok_states:
|
||||
raise ValueError("txn in state %r but expected one of %r" %
|
||||
(self.txn_state, ok_states))
|
||||
|
||||
def _checkTransaction(self, transaction):
|
||||
if (transaction is not self.transaction
|
||||
and self.transaction is not None):
|
||||
raise TypeError("Transaction missmatch",
|
||||
transaction, self.transaction)
|
||||
|
||||
def inc(self, n=1):
|
||||
self.delta += n
|
||||
|
||||
def tpc_begin(self, transaction):
|
||||
self._checkTransaction(transaction)
|
||||
self._check_state(None)
|
||||
self.transaction = transaction
|
||||
self.txn_state = 'tpc_begin'
|
||||
|
||||
def tpc_vote(self, transaction):
|
||||
self._checkTransaction(transaction)
|
||||
self._check_state('tpc_begin')
|
||||
self.state += self.delta
|
||||
self.txn_state = 'tpc_vote'
|
||||
|
||||
def tpc_finish(self, transaction):
|
||||
self._checkTransaction(transaction)
|
||||
self._check_state('tpc_vote')
|
||||
self.delta = 0
|
||||
self.transaction = None
|
||||
self.prepared = False
|
||||
self.txn_state = None
|
||||
|
||||
def tpc_abort(self, transaction):
|
||||
self._checkTransaction(transaction)
|
||||
if self.transaction is not None:
|
||||
self.transaction = None
|
||||
|
||||
if self.txn_state == 'tpc_vote':
|
||||
self.state -= self.delta
|
||||
|
||||
self.txn_state = None
|
||||
self.delta = 0
|
||||
|
||||
def savepoint(self, transaction):
|
||||
if self.txn_state is not None:
|
||||
raise AssertionError("Can't get savepoint during two-phase commit")
|
||||
self._checkTransaction(transaction)
|
||||
self.transaction = transaction
|
||||
self.sp += 1
|
||||
return SavePoint(self)
|
||||
|
||||
def discard(self, transaction):
|
||||
"Does nothing"
|
||||
|
||||
|
||||
class SavePoint(object):
|
||||
|
||||
def __init__(self, rm):
|
||||
self.rm = rm
|
||||
self.sp = rm.sp
|
||||
self.delta = rm.delta
|
||||
self.transaction = rm.transaction
|
||||
|
||||
def rollback(self):
|
||||
if self.transaction is not self.rm.transaction:
|
||||
raise TypeError("Attempt to rollback stale rollback")
|
||||
if self.rm.sp < self.sp:
|
||||
raise TypeError("Attempt to roll back to invalid save point",
|
||||
self.sp, self.rm.sp)
|
||||
self.rm.sp = self.sp
|
||||
self.rm.delta = self.delta
|
||||
|
||||
def discard(self):
|
||||
"Does nothing."
|
||||
@@ -0,0 +1,185 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
##############################################################################
|
||||
"""Savepoint data manager implementation example.
|
||||
|
||||
Sample data manager implementation that illustrates how to implement
|
||||
savepoints.
|
||||
|
||||
Used by savepoint.rst in the Sphinx docs.
|
||||
"""
|
||||
|
||||
from zope.interface import implementer
|
||||
import transaction.interfaces
|
||||
|
||||
@implementer(transaction.interfaces.IDataManager)
|
||||
class SampleDataManager(object):
|
||||
"""Sample implementation of data manager that doesn't support savepoints
|
||||
|
||||
This data manager stores named simple values, like strings and numbers.
|
||||
"""
|
||||
|
||||
def __init__(self, transaction_manager=None):
|
||||
if transaction_manager is None:
|
||||
# Use the thread-local transaction manager if none is provided:
|
||||
import transaction
|
||||
transaction_manager = transaction.manager
|
||||
self.transaction_manager = transaction_manager
|
||||
|
||||
# Our committed and uncommitted data:
|
||||
self.committed = {}
|
||||
self.uncommitted = self.committed.copy()
|
||||
|
||||
# Our transaction state:
|
||||
#
|
||||
# If our uncommitted data is modified, we'll join a transaction
|
||||
# and keep track of the transaction we joined. Any commit
|
||||
# related messages we get should be for this same transaction
|
||||
self.transaction = None
|
||||
|
||||
# What phase, if any, of two-phase commit we are in:
|
||||
self.tpc_phase = None
|
||||
|
||||
|
||||
#######################################################################
|
||||
# Provide a mapping interface to uncommitted data. We provide
|
||||
# a basic subset of the interface. DictMixin does the rest.
|
||||
|
||||
def __getitem__(self, name):
|
||||
return self.uncommitted[name]
|
||||
|
||||
def __setitem__(self, name, value):
|
||||
self._join() # join the current transaction, if we haven't already
|
||||
self.uncommitted[name] = value
|
||||
|
||||
def keys(self):
|
||||
return self.uncommitted.keys()
|
||||
|
||||
__iter__ = keys
|
||||
|
||||
def __contains__(self, k):
|
||||
return k in self.uncommitted
|
||||
|
||||
def __repr__(self):
|
||||
return repr(self.uncommitted)
|
||||
|
||||
#
|
||||
#######################################################################
|
||||
|
||||
#######################################################################
|
||||
# Transaction methods
|
||||
|
||||
def _join(self):
|
||||
# If this is the first change in the transaction, join the transaction
|
||||
if self.transaction is None:
|
||||
self.transaction = self.transaction_manager.get()
|
||||
self.transaction.join(self)
|
||||
|
||||
def _resetTransaction(self):
|
||||
self.last_note = getattr(self.transaction, 'description', None)
|
||||
self.transaction = None
|
||||
self.tpc_phase = None
|
||||
|
||||
def abort(self, transaction):
|
||||
"""Throw away changes made before the commit process has started
|
||||
"""
|
||||
assert ((transaction is self.transaction) or (self.transaction is None)
|
||||
), "Must not change transactions"
|
||||
assert self.tpc_phase is None, "Must be called outside of tpc"
|
||||
self.uncommitted = self.committed.copy()
|
||||
self._resetTransaction()
|
||||
|
||||
def tpc_begin(self, transaction):
|
||||
"""Enter two-phase commit
|
||||
"""
|
||||
assert transaction is self.transaction, "Must not change transactions"
|
||||
assert self.tpc_phase is None, "Must be called outside of tpc"
|
||||
self.tpc_phase = 1
|
||||
|
||||
def commit(self, transaction):
|
||||
"""Record data modified during the transaction
|
||||
"""
|
||||
assert transaction is self.transaction, "Must not change transactions"
|
||||
assert self.tpc_phase == 1, "Must be called in first phase of tpc"
|
||||
|
||||
# In our simple example, we don't need to do anything.
|
||||
# A more complex data manager would typically write to some sort
|
||||
# of log.
|
||||
|
||||
def tpc_vote(self, transaction):
|
||||
assert transaction is self.transaction, "Must not change transactions"
|
||||
assert self.tpc_phase == 1, "Must be called in first phase of tpc"
|
||||
# This particular data manager is always ready to vote.
|
||||
# Real data managers will usually need to take some steps to
|
||||
# make sure that the finish will succeed
|
||||
self.tpc_phase = 2
|
||||
|
||||
def tpc_finish(self, transaction):
|
||||
assert transaction is self.transaction, "Must not change transactions"
|
||||
assert self.tpc_phase == 2, "Must be called in second phase of tpc"
|
||||
self.committed = self.uncommitted.copy()
|
||||
self._resetTransaction()
|
||||
|
||||
def tpc_abort(self, transaction):
|
||||
if self.transaction is not None: # pragma: no cover
|
||||
# otherwise we're not actually joined.
|
||||
assert self.tpc_phase is not None, "Must be called inside of tpc"
|
||||
self.uncommitted = self.committed.copy()
|
||||
self._resetTransaction()
|
||||
|
||||
#
|
||||
#######################################################################
|
||||
|
||||
#######################################################################
|
||||
# Other data manager methods
|
||||
|
||||
def sortKey(self):
|
||||
# Commit operations on multiple data managers are performed in
|
||||
# sort key order. This important to avoid deadlock when data
|
||||
# managers are shared among multiple threads or processes and
|
||||
# use locks to manage that sharing. We aren't going to bother
|
||||
# with that here.
|
||||
return str(id(self))
|
||||
|
||||
#
|
||||
#######################################################################
|
||||
|
||||
@implementer(transaction.interfaces.ISavepointDataManager)
|
||||
class SampleSavepointDataManager(SampleDataManager):
|
||||
"""Sample implementation of a savepoint-supporting data manager
|
||||
|
||||
This extends the basic data manager with savepoint support.
|
||||
"""
|
||||
|
||||
def savepoint(self):
|
||||
# When we create the savepoint, we save the existing database state.
|
||||
return SampleSavepoint(self, self.uncommitted.copy())
|
||||
|
||||
def _rollback_savepoint(self, savepoint):
|
||||
# When we rollback the savepoint, we restore the saved data.
|
||||
# Caution: without the copy(), further changes to the database
|
||||
# could reflect in savepoint.data, and then `savepoint` would no
|
||||
# longer contain the originally saved data, and so `savepoint`
|
||||
# couldn't restore the original state if a rollback to this
|
||||
# savepoint was done again. IOW, copy() is necessary.
|
||||
self.uncommitted = savepoint.data.copy()
|
||||
|
||||
@implementer(transaction.interfaces.IDataManagerSavepoint)
|
||||
class SampleSavepoint:
|
||||
|
||||
def __init__(self, data_manager, data):
|
||||
self.data_manager = data_manager
|
||||
self.data = data
|
||||
|
||||
def rollback(self):
|
||||
self.data_manager._rollback_savepoint(self)
|
||||
@@ -0,0 +1,991 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2012 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE
|
||||
#
|
||||
##############################################################################
|
||||
import mock
|
||||
import unittest
|
||||
|
||||
import zope.interface.verify
|
||||
|
||||
from .. import interfaces
|
||||
|
||||
|
||||
class TransactionManagerTests(unittest.TestCase):
|
||||
|
||||
def _getTargetClass(self):
|
||||
from transaction import TransactionManager
|
||||
return TransactionManager
|
||||
|
||||
def _makeOne(self):
|
||||
return self._getTargetClass()()
|
||||
|
||||
def _makePopulated(self):
|
||||
mgr = self._makeOne()
|
||||
sub1 = DataObject(mgr)
|
||||
sub2 = DataObject(mgr)
|
||||
sub3 = DataObject(mgr)
|
||||
nosub1 = DataObject(mgr, nost=1)
|
||||
return mgr, sub1, sub2, sub3, nosub1
|
||||
|
||||
def test_interface(self):
|
||||
zope.interface.verify.verifyObject(interfaces.ITransactionManager,
|
||||
self._makeOne())
|
||||
|
||||
def test_ctor(self):
|
||||
tm = self._makeOne()
|
||||
self.assertTrue(tm._txn is None)
|
||||
self.assertEqual(len(tm._synchs), 0)
|
||||
|
||||
def test_begin_wo_existing_txn_wo_synchs(self):
|
||||
from transaction._transaction import Transaction
|
||||
tm = self._makeOne()
|
||||
tm.begin()
|
||||
self.assertTrue(isinstance(tm._txn, Transaction))
|
||||
|
||||
def test_begin_wo_existing_txn_w_synchs(self):
|
||||
from transaction._transaction import Transaction
|
||||
tm = self._makeOne()
|
||||
synch = DummySynch()
|
||||
tm.registerSynch(synch)
|
||||
tm.begin()
|
||||
self.assertTrue(isinstance(tm._txn, Transaction))
|
||||
self.assertTrue(tm._txn in synch._txns)
|
||||
|
||||
def test_begin_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_aborted = False
|
||||
def abort(self):
|
||||
self._aborted = True
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
tm.begin()
|
||||
self.assertFalse(tm._txn is txn)
|
||||
self.assertTrue(txn._aborted)
|
||||
|
||||
def test_get_wo_existing_txn(self):
|
||||
from transaction._transaction import Transaction
|
||||
tm = self._makeOne()
|
||||
txn = tm.get()
|
||||
self.assertTrue(isinstance(txn, Transaction))
|
||||
|
||||
def test_get_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_aborted = False
|
||||
def abort(self):
|
||||
raise AssertionError("This is not actually called")
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
self.assertIs(tm.get(), txn)
|
||||
|
||||
def test_free_w_other_txn(self):
|
||||
from transaction._transaction import Transaction
|
||||
tm = self._makeOne()
|
||||
txn = Transaction()
|
||||
tm.begin()
|
||||
self.assertRaises(ValueError, tm.free, txn)
|
||||
|
||||
def test_free_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_aborted = False
|
||||
def abort(self):
|
||||
raise AssertionError("This is not actually called")
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
tm.free(txn)
|
||||
self.assertIsNone(tm._txn)
|
||||
|
||||
def test_registerSynch(self):
|
||||
tm = self._makeOne()
|
||||
synch = DummySynch()
|
||||
tm.registerSynch(synch)
|
||||
self.assertEqual(len(tm._synchs), 1)
|
||||
self.assertTrue(synch in tm._synchs)
|
||||
|
||||
def test_unregisterSynch(self):
|
||||
tm = self._makeOne()
|
||||
synch1 = DummySynch()
|
||||
synch2 = DummySynch()
|
||||
self.assertFalse(tm.registeredSynchs())
|
||||
tm.registerSynch(synch1)
|
||||
self.assertTrue(tm.registeredSynchs())
|
||||
tm.registerSynch(synch2)
|
||||
self.assertTrue(tm.registeredSynchs())
|
||||
tm.unregisterSynch(synch1)
|
||||
self.assertTrue(tm.registeredSynchs())
|
||||
self.assertEqual(len(tm._synchs), 1)
|
||||
self.assertFalse(synch1 in tm._synchs)
|
||||
self.assertTrue(synch2 in tm._synchs)
|
||||
tm.unregisterSynch(synch2)
|
||||
self.assertFalse(tm.registeredSynchs())
|
||||
|
||||
def test_clearSynchs(self):
|
||||
tm = self._makeOne()
|
||||
synch1 = DummySynch()
|
||||
synch2 = DummySynch()
|
||||
tm.registerSynch(synch1)
|
||||
tm.registerSynch(synch2)
|
||||
tm.clearSynchs()
|
||||
self.assertEqual(len(tm._synchs), 0)
|
||||
|
||||
def test_isDoomed_wo_existing_txn(self):
|
||||
tm = self._makeOne()
|
||||
self.assertFalse(tm.isDoomed())
|
||||
tm._txn.doom()
|
||||
self.assertTrue(tm.isDoomed())
|
||||
|
||||
def test_isDoomed_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_doomed = False
|
||||
def isDoomed(self):
|
||||
return self._doomed
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
self.assertFalse(tm.isDoomed())
|
||||
txn._doomed = True
|
||||
self.assertTrue(tm.isDoomed())
|
||||
|
||||
def test_doom(self):
|
||||
tm = self._makeOne()
|
||||
txn = tm.get()
|
||||
self.assertFalse(txn.isDoomed())
|
||||
tm.doom()
|
||||
self.assertTrue(txn.isDoomed())
|
||||
self.assertTrue(tm.isDoomed())
|
||||
|
||||
def test_commit_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_committed = False
|
||||
def commit(self):
|
||||
self._committed = True
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
tm.commit()
|
||||
self.assertTrue(txn._committed)
|
||||
|
||||
def test_abort_w_existing_txn(self):
|
||||
class Existing(object):
|
||||
_aborted = False
|
||||
def abort(self):
|
||||
self._aborted = True
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = Existing()
|
||||
tm.abort()
|
||||
self.assertTrue(txn._aborted)
|
||||
|
||||
def test_as_context_manager_wo_error(self):
|
||||
class _Test(object):
|
||||
_committed = False
|
||||
_aborted = False
|
||||
def commit(self):
|
||||
self._committed = True
|
||||
def abort(self):
|
||||
raise AssertionError("This should not be called")
|
||||
tm = self._makeOne()
|
||||
with tm:
|
||||
tm._txn = txn = _Test()
|
||||
self.assertTrue(txn._committed)
|
||||
self.assertFalse(txn._aborted)
|
||||
|
||||
def test_as_context_manager_w_error(self):
|
||||
class _Test(object):
|
||||
_committed = False
|
||||
_aborted = False
|
||||
def commit(self):
|
||||
raise AssertionError("This should not be called")
|
||||
def abort(self):
|
||||
self._aborted = True
|
||||
tm = self._makeOne()
|
||||
|
||||
with self.assertRaises(ZeroDivisionError):
|
||||
with tm:
|
||||
tm._txn = txn = _Test()
|
||||
raise ZeroDivisionError()
|
||||
|
||||
self.assertFalse(txn._committed)
|
||||
self.assertTrue(txn._aborted)
|
||||
|
||||
def test_savepoint_default(self):
|
||||
class _Test(object):
|
||||
_sp = None
|
||||
def savepoint(self, optimistic):
|
||||
self._sp = optimistic
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = _Test()
|
||||
tm.savepoint()
|
||||
self.assertFalse(txn._sp)
|
||||
|
||||
def test_savepoint_explicit(self):
|
||||
class _Test(object):
|
||||
_sp = None
|
||||
def savepoint(self, optimistic):
|
||||
self._sp = optimistic
|
||||
tm = self._makeOne()
|
||||
tm._txn = txn = _Test()
|
||||
tm.savepoint(True)
|
||||
self.assertTrue(txn._sp)
|
||||
|
||||
def test_attempts_w_invalid_count(self):
|
||||
tm = self._makeOne()
|
||||
self.assertRaises(ValueError, list, tm.attempts(0))
|
||||
self.assertRaises(ValueError, list, tm.attempts(-1))
|
||||
self.assertRaises(ValueError, list, tm.attempts(-10))
|
||||
|
||||
def test_attempts_w_valid_count(self):
|
||||
tm = self._makeOne()
|
||||
found = list(tm.attempts(1))
|
||||
self.assertEqual(len(found), 1)
|
||||
self.assertTrue(found[0] is tm)
|
||||
|
||||
def test_attempts_stop_on_success(self):
|
||||
tm = self._makeOne()
|
||||
|
||||
i = 0
|
||||
for attempt in tm.attempts():
|
||||
with attempt:
|
||||
i += 1
|
||||
|
||||
self.assertEqual(i, 1)
|
||||
|
||||
def test_attempts_retries(self):
|
||||
import transaction.interfaces
|
||||
class Retry(transaction.interfaces.TransientError):
|
||||
pass
|
||||
|
||||
tm = self._makeOne()
|
||||
i = 0
|
||||
for attempt in tm.attempts(4):
|
||||
with attempt:
|
||||
i += 1
|
||||
if i < 4:
|
||||
raise Retry
|
||||
|
||||
self.assertEqual(i, 4)
|
||||
|
||||
def test_attempts_retries_but_gives_up(self):
|
||||
import transaction.interfaces
|
||||
class Retry(transaction.interfaces.TransientError):
|
||||
pass
|
||||
|
||||
tm = self._makeOne()
|
||||
i = 0
|
||||
|
||||
with self.assertRaises(Retry):
|
||||
for attempt in tm.attempts(4):
|
||||
with attempt:
|
||||
i += 1
|
||||
raise Retry
|
||||
|
||||
self.assertEqual(i, 4)
|
||||
|
||||
def test_attempts_propigates_errors(self):
|
||||
tm = self._makeOne()
|
||||
with self.assertRaises(ValueError):
|
||||
for attempt in tm.attempts(4):
|
||||
with attempt:
|
||||
raise ValueError
|
||||
|
||||
def test_attempts_defer_to_dm(self):
|
||||
import transaction.tests.savepointsample
|
||||
|
||||
class DM(transaction.tests.savepointsample.SampleSavepointDataManager):
|
||||
def should_retry(self, e):
|
||||
if 'should retry' in str(e):
|
||||
return True
|
||||
|
||||
ntry = 0
|
||||
dm = transaction.tests.savepointsample.SampleSavepointDataManager()
|
||||
dm2 = DM()
|
||||
with transaction.manager:
|
||||
dm2['ntry'] = 0
|
||||
|
||||
for attempt in transaction.manager.attempts():
|
||||
with attempt:
|
||||
ntry += 1
|
||||
dm['ntry'] = ntry
|
||||
dm2['ntry'] = ntry
|
||||
if ntry % 3:
|
||||
raise ValueError('we really should retry this')
|
||||
|
||||
self.assertEqual(ntry, 3)
|
||||
|
||||
|
||||
def test_attempts_w_default_count(self):
|
||||
from transaction._manager import Attempt
|
||||
tm = self._makeOne()
|
||||
found = list(tm.attempts())
|
||||
self.assertEqual(len(found), 3)
|
||||
for attempt in found[:-1]:
|
||||
self.assertTrue(isinstance(attempt, Attempt))
|
||||
self.assertTrue(attempt.manager is tm)
|
||||
self.assertTrue(found[-1] is tm)
|
||||
|
||||
def test_run(self):
|
||||
import transaction.interfaces
|
||||
class Retry(transaction.interfaces.TransientError):
|
||||
pass
|
||||
|
||||
tm = self._makeOne()
|
||||
i = [0, None]
|
||||
|
||||
@tm.run()
|
||||
def meaning():
|
||||
"Nice doc"
|
||||
i[0] += 1
|
||||
i[1] = tm.get()
|
||||
if i[0] < 3:
|
||||
raise Retry
|
||||
return 42
|
||||
|
||||
self.assertEqual(i[0], 3)
|
||||
self.assertEqual(meaning, 42)
|
||||
self.assertEqual(i[1].description, "meaning\n\nNice doc")
|
||||
|
||||
def test_run_no_name_explicit_tries(self):
|
||||
import transaction.interfaces
|
||||
class Retry(transaction.interfaces.TransientError):
|
||||
pass
|
||||
|
||||
tm = self._makeOne()
|
||||
i = [0, None]
|
||||
|
||||
@tm.run(4)
|
||||
def _():
|
||||
"Nice doc"
|
||||
i[0] += 1
|
||||
i[1] = tm.get()
|
||||
if i[0] < 4:
|
||||
raise Retry
|
||||
|
||||
self.assertEqual(i[0], 4)
|
||||
self.assertEqual(i[1].description, "Nice doc")
|
||||
|
||||
def test_run_pos_tries(self):
|
||||
tm = self._makeOne()
|
||||
|
||||
with self.assertRaises(ValueError):
|
||||
tm.run(0)(lambda : None)
|
||||
with self.assertRaises(ValueError):
|
||||
@tm.run(-1)
|
||||
def _():
|
||||
raise AssertionError("Never called")
|
||||
|
||||
def test_run_stop_on_success(self):
|
||||
import transaction.interfaces
|
||||
|
||||
tm = self._makeOne()
|
||||
i = [0, None]
|
||||
|
||||
@tm.run()
|
||||
def meaning():
|
||||
i[0] += 1
|
||||
i[1] = tm.get()
|
||||
return 43
|
||||
|
||||
self.assertEqual(i[0], 1)
|
||||
self.assertEqual(meaning, 43)
|
||||
self.assertEqual(i[1].description, "meaning")
|
||||
|
||||
def test_run_retries_but_gives_up(self):
|
||||
import transaction.interfaces
|
||||
class Retry(transaction.interfaces.TransientError):
|
||||
pass
|
||||
|
||||
tm = self._makeOne()
|
||||
i = [0]
|
||||
|
||||
with self.assertRaises(Retry):
|
||||
@tm.run()
|
||||
def _():
|
||||
i[0] += 1
|
||||
raise Retry
|
||||
|
||||
self.assertEqual(i[0], 3)
|
||||
|
||||
def test_run_propigates_errors(self):
|
||||
tm = self._makeOne()
|
||||
with self.assertRaises(ValueError):
|
||||
@tm.run
|
||||
def _():
|
||||
raise ValueError
|
||||
|
||||
def test_run_defer_to_dm(self):
|
||||
import transaction.tests.savepointsample
|
||||
|
||||
class DM(transaction.tests.savepointsample.SampleSavepointDataManager):
|
||||
def should_retry(self, e):
|
||||
if 'should retry' in str(e):
|
||||
return True
|
||||
|
||||
ntry = [0]
|
||||
dm = transaction.tests.savepointsample.SampleSavepointDataManager()
|
||||
dm2 = DM()
|
||||
with transaction.manager:
|
||||
dm2['ntry'] = 0
|
||||
|
||||
@transaction.manager.run
|
||||
def _():
|
||||
ntry[0] += 1
|
||||
dm['ntry'] = ntry[0]
|
||||
dm2['ntry'] = ntry[0]
|
||||
if ntry[0] % 3:
|
||||
raise ValueError('we really should retry this')
|
||||
|
||||
self.assertEqual(ntry[0], 3)
|
||||
|
||||
def test_run_callable_with_bytes_doc(self):
|
||||
import transaction
|
||||
class Callable(object):
|
||||
|
||||
def __init__(self):
|
||||
self.__doc__ = b'some bytes'
|
||||
self.__name__ = b'more bytes'
|
||||
|
||||
def __call__(self):
|
||||
return 42
|
||||
|
||||
result = transaction.manager.run(Callable())
|
||||
self.assertEqual(result, 42)
|
||||
|
||||
def test__retryable_w_transient_error(self):
|
||||
from transaction.interfaces import TransientError
|
||||
tm = self._makeOne()
|
||||
self.assertTrue(tm._retryable(TransientError, object()))
|
||||
|
||||
def test__retryable_w_transient_subclass(self):
|
||||
from transaction.interfaces import TransientError
|
||||
class _Derived(TransientError):
|
||||
pass
|
||||
tm = self._makeOne()
|
||||
self.assertTrue(tm._retryable(_Derived, object()))
|
||||
|
||||
def test__retryable_w_normal_exception_no_resources(self):
|
||||
tm = self._makeOne()
|
||||
self.assertFalse(tm._retryable(Exception, object()))
|
||||
|
||||
def test__retryable_w_normal_exception_w_resource_voting_yes(self):
|
||||
class _Resource(object):
|
||||
def should_retry(self, err):
|
||||
return True
|
||||
tm = self._makeOne()
|
||||
tm.get()._resources.append(_Resource())
|
||||
self.assertTrue(tm._retryable(Exception, object()))
|
||||
|
||||
def test__retryable_w_multiple(self):
|
||||
class _Resource(object):
|
||||
_should = True
|
||||
def should_retry(self, err):
|
||||
return self._should
|
||||
tm = self._makeOne()
|
||||
res1 = _Resource()
|
||||
res1._should = False
|
||||
res2 = _Resource()
|
||||
tm.get()._resources.append(res1)
|
||||
tm.get()._resources.append(res2)
|
||||
self.assertTrue(tm._retryable(Exception, object()))
|
||||
|
||||
# basic tests with two sub trans jars
|
||||
# really we only need one, so tests for
|
||||
# sub1 should identical to tests for sub2
|
||||
def test_commit_normal(self):
|
||||
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1.modify()
|
||||
sub2.modify()
|
||||
|
||||
mgr.commit()
|
||||
|
||||
assert sub1._p_jar.ccommit_sub == 0
|
||||
assert sub1._p_jar.ctpc_finish == 1
|
||||
|
||||
def test_abort_normal(self):
|
||||
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1.modify()
|
||||
sub2.modify()
|
||||
|
||||
mgr.abort()
|
||||
|
||||
assert sub2._p_jar.cabort == 1
|
||||
|
||||
|
||||
# repeat adding in a nonsub trans jars
|
||||
|
||||
def test_commit_w_nonsub_jar(self):
|
||||
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
nosub1.modify()
|
||||
|
||||
mgr.commit()
|
||||
|
||||
assert nosub1._p_jar.ctpc_finish == 1
|
||||
|
||||
def test_abort_w_nonsub_jar(self):
|
||||
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
nosub1.modify()
|
||||
|
||||
mgr.abort()
|
||||
|
||||
assert nosub1._p_jar.ctpc_finish == 0
|
||||
assert nosub1._p_jar.cabort == 1
|
||||
|
||||
|
||||
### Failure Mode Tests
|
||||
#
|
||||
# ok now we do some more interesting
|
||||
# tests that check the implementations
|
||||
# error handling by throwing errors from
|
||||
# various jar methods
|
||||
###
|
||||
|
||||
# first the recoverable errors
|
||||
|
||||
def test_abort_w_broken_jar(self):
|
||||
from transaction import _transaction
|
||||
from transaction.tests.common import DummyLogger
|
||||
from transaction.tests.common import Monkey
|
||||
logger = DummyLogger()
|
||||
with Monkey(_transaction, _LOGGER=logger):
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1._p_jar = BasicJar(errors='abort')
|
||||
nosub1.modify()
|
||||
sub1.modify(nojar=1)
|
||||
sub2.modify()
|
||||
try:
|
||||
mgr.abort()
|
||||
except TestTxnException:
|
||||
pass
|
||||
|
||||
assert nosub1._p_jar.cabort == 1
|
||||
assert sub2._p_jar.cabort == 1
|
||||
|
||||
def test_commit_w_broken_jar_commit(self):
|
||||
from transaction import _transaction
|
||||
from transaction.tests.common import DummyLogger
|
||||
from transaction.tests.common import Monkey
|
||||
logger = DummyLogger()
|
||||
with Monkey(_transaction, _LOGGER=logger):
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1._p_jar = BasicJar(errors='commit')
|
||||
nosub1.modify()
|
||||
sub1.modify(nojar=1)
|
||||
try:
|
||||
mgr.commit()
|
||||
except TestTxnException:
|
||||
pass
|
||||
|
||||
assert nosub1._p_jar.ctpc_finish == 0
|
||||
assert nosub1._p_jar.ccommit == 1
|
||||
assert nosub1._p_jar.ctpc_abort == 1
|
||||
|
||||
def test_commit_w_broken_jar_tpc_vote(self):
|
||||
from transaction import _transaction
|
||||
from transaction.tests.common import DummyLogger
|
||||
from transaction.tests.common import Monkey
|
||||
logger = DummyLogger()
|
||||
with Monkey(_transaction, _LOGGER=logger):
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1._p_jar = BasicJar(errors='tpc_vote')
|
||||
nosub1.modify()
|
||||
sub1.modify(nojar=1)
|
||||
try:
|
||||
mgr.commit()
|
||||
except TestTxnException:
|
||||
pass
|
||||
|
||||
assert nosub1._p_jar.ctpc_finish == 0
|
||||
assert nosub1._p_jar.ccommit == 1
|
||||
assert nosub1._p_jar.ctpc_abort == 1
|
||||
assert sub1._p_jar.ctpc_abort == 1
|
||||
|
||||
def test_commit_w_broken_jar_tpc_begin(self):
|
||||
# ok this test reveals a bug in the TM.py
|
||||
# as the nosub tpc_abort there is ignored.
|
||||
|
||||
# nosub calling method tpc_begin
|
||||
# nosub calling method commit
|
||||
# sub calling method tpc_begin
|
||||
# sub calling method abort
|
||||
# sub calling method tpc_abort
|
||||
# nosub calling method tpc_abort
|
||||
from transaction import _transaction
|
||||
from transaction.tests.common import DummyLogger
|
||||
from transaction.tests.common import Monkey
|
||||
logger = DummyLogger()
|
||||
with Monkey(_transaction, _LOGGER=logger):
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1._p_jar = BasicJar(errors='tpc_begin')
|
||||
nosub1.modify()
|
||||
sub1.modify(nojar=1)
|
||||
try:
|
||||
mgr.commit()
|
||||
except TestTxnException:
|
||||
pass
|
||||
|
||||
assert nosub1._p_jar.ctpc_abort == 1
|
||||
assert sub1._p_jar.ctpc_abort == 1
|
||||
|
||||
def test_commit_w_broken_jar_tpc_abort_tpc_vote(self):
|
||||
from transaction import _transaction
|
||||
from transaction.tests.common import DummyLogger
|
||||
from transaction.tests.common import Monkey
|
||||
logger = DummyLogger()
|
||||
with Monkey(_transaction, _LOGGER=logger):
|
||||
mgr, sub1, sub2, sub3, nosub1 = self._makePopulated()
|
||||
sub1._p_jar = BasicJar(errors=('tpc_abort', 'tpc_vote'))
|
||||
nosub1.modify()
|
||||
sub1.modify(nojar=1)
|
||||
try:
|
||||
mgr.commit()
|
||||
except TestTxnException:
|
||||
pass
|
||||
|
||||
assert nosub1._p_jar.ctpc_abort == 1
|
||||
|
||||
def test_notify_transaction_late_comers(self):
|
||||
# If a datamanager registers for synchonization after a
|
||||
# transaction has started, we should call newTransaction so it
|
||||
# can do necessry setup.
|
||||
import mock
|
||||
from .. import TransactionManager
|
||||
manager = TransactionManager()
|
||||
sync1 = mock.MagicMock()
|
||||
manager.registerSynch(sync1)
|
||||
sync1.newTransaction.assert_not_called()
|
||||
t = manager.begin()
|
||||
sync1.newTransaction.assert_called_with(t)
|
||||
sync2 = mock.MagicMock()
|
||||
manager.registerSynch(sync2)
|
||||
sync2.newTransaction.assert_called_with(t)
|
||||
|
||||
# for, um, completeness
|
||||
t.commit()
|
||||
for s in sync1, sync2:
|
||||
s.beforeCompletion.assert_called_with(t)
|
||||
s.afterCompletion.assert_called_with(t)
|
||||
|
||||
def test_unregisterSynch_on_transaction_manager_from_serparate_thread(self):
|
||||
# We should be able to get the underlying manager of the thread manager
|
||||
# cand call methods from other threads.
|
||||
|
||||
import threading, transaction
|
||||
|
||||
started = threading.Event()
|
||||
stopped = threading.Event()
|
||||
|
||||
synchronizer = self
|
||||
|
||||
class Runner(threading.Thread):
|
||||
|
||||
def __init__(self):
|
||||
threading.Thread.__init__(self)
|
||||
self.manager = transaction.manager.manager
|
||||
self.setDaemon(True)
|
||||
self.start()
|
||||
|
||||
def run(self):
|
||||
self.manager.registerSynch(synchronizer)
|
||||
started.set()
|
||||
stopped.wait()
|
||||
|
||||
runner = Runner()
|
||||
started.wait()
|
||||
runner.manager.unregisterSynch(synchronizer)
|
||||
stopped.set()
|
||||
runner.join(1)
|
||||
|
||||
|
||||
class TestThreadTransactionManager(unittest.TestCase):
|
||||
|
||||
def test_interface(self):
|
||||
import transaction
|
||||
zope.interface.verify.verifyObject(interfaces.ITransactionManager,
|
||||
transaction.manager)
|
||||
|
||||
def test_sync_registration_thread_local_manager(self):
|
||||
import transaction
|
||||
|
||||
sync = mock.MagicMock()
|
||||
sync2 = mock.MagicMock()
|
||||
self.assertFalse(transaction.manager.registeredSynchs())
|
||||
transaction.manager.registerSynch(sync)
|
||||
self.assertTrue(transaction.manager.registeredSynchs())
|
||||
transaction.manager.registerSynch(sync2)
|
||||
self.assertTrue(transaction.manager.registeredSynchs())
|
||||
t = transaction.begin()
|
||||
sync.newTransaction.assert_called_with(t)
|
||||
transaction.abort()
|
||||
sync.beforeCompletion.assert_called_with(t)
|
||||
sync.afterCompletion.assert_called_with(t)
|
||||
transaction.manager.unregisterSynch(sync)
|
||||
self.assertTrue(transaction.manager.registeredSynchs())
|
||||
transaction.manager.unregisterSynch(sync2)
|
||||
self.assertFalse(transaction.manager.registeredSynchs())
|
||||
sync.reset_mock()
|
||||
transaction.begin()
|
||||
transaction.abort()
|
||||
sync.newTransaction.assert_not_called()
|
||||
sync.beforeCompletion.assert_not_called()
|
||||
sync.afterCompletion.assert_not_called()
|
||||
|
||||
self.assertFalse(transaction.manager.registeredSynchs())
|
||||
transaction.manager.registerSynch(sync)
|
||||
transaction.manager.registerSynch(sync2)
|
||||
t = transaction.begin()
|
||||
sync.newTransaction.assert_called_with(t)
|
||||
self.assertTrue(transaction.manager.registeredSynchs())
|
||||
transaction.abort()
|
||||
sync.beforeCompletion.assert_called_with(t)
|
||||
sync.afterCompletion.assert_called_with(t)
|
||||
transaction.manager.clearSynchs()
|
||||
self.assertFalse(transaction.manager.registeredSynchs())
|
||||
sync.reset_mock()
|
||||
transaction.begin()
|
||||
transaction.abort()
|
||||
sync.newTransaction.assert_not_called()
|
||||
sync.beforeCompletion.assert_not_called()
|
||||
sync.afterCompletion.assert_not_called()
|
||||
|
||||
def test_explicit_thread_local_manager(self):
|
||||
import transaction.interfaces
|
||||
|
||||
self.assertFalse(transaction.manager.explicit)
|
||||
transaction.abort()
|
||||
transaction.manager.explicit = True
|
||||
self.assertTrue(transaction.manager.explicit)
|
||||
with self.assertRaises(transaction.interfaces.NoTransaction):
|
||||
transaction.abort()
|
||||
transaction.manager.explicit = False
|
||||
transaction.abort()
|
||||
|
||||
|
||||
class AttemptTests(unittest.TestCase):
|
||||
|
||||
def _makeOne(self, manager):
|
||||
from transaction._manager import Attempt
|
||||
return Attempt(manager)
|
||||
|
||||
def test___enter__(self):
|
||||
manager = DummyManager()
|
||||
inst = self._makeOne(manager)
|
||||
inst.__enter__()
|
||||
self.assertTrue(manager.entered)
|
||||
|
||||
def test___exit__no_exc_no_commit_exception(self):
|
||||
manager = DummyManager()
|
||||
inst = self._makeOne(manager)
|
||||
result = inst.__exit__(None, None, None)
|
||||
self.assertFalse(result)
|
||||
self.assertTrue(manager.committed)
|
||||
|
||||
def test___exit__no_exc_nonretryable_commit_exception(self):
|
||||
manager = DummyManager(raise_on_commit=ValueError)
|
||||
inst = self._makeOne(manager)
|
||||
self.assertRaises(ValueError, inst.__exit__, None, None, None)
|
||||
self.assertTrue(manager.committed)
|
||||
self.assertTrue(manager.aborted)
|
||||
|
||||
def test___exit__no_exc_abort_exception_after_nonretryable_commit_exc(self):
|
||||
manager = DummyManager(raise_on_abort=ValueError,
|
||||
raise_on_commit=KeyError)
|
||||
inst = self._makeOne(manager)
|
||||
self.assertRaises(ValueError, inst.__exit__, None, None, None)
|
||||
self.assertTrue(manager.committed)
|
||||
self.assertTrue(manager.aborted)
|
||||
|
||||
def test___exit__no_exc_retryable_commit_exception(self):
|
||||
from transaction.interfaces import TransientError
|
||||
manager = DummyManager(raise_on_commit=TransientError)
|
||||
inst = self._makeOne(manager)
|
||||
result = inst.__exit__(None, None, None)
|
||||
self.assertTrue(result)
|
||||
self.assertTrue(manager.committed)
|
||||
self.assertTrue(manager.aborted)
|
||||
|
||||
def test___exit__with_exception_value_retryable(self):
|
||||
from transaction.interfaces import TransientError
|
||||
manager = DummyManager()
|
||||
inst = self._makeOne(manager)
|
||||
result = inst.__exit__(TransientError, TransientError(), None)
|
||||
self.assertTrue(result)
|
||||
self.assertFalse(manager.committed)
|
||||
self.assertTrue(manager.aborted)
|
||||
|
||||
def test___exit__with_exception_value_nonretryable(self):
|
||||
manager = DummyManager()
|
||||
inst = self._makeOne(manager)
|
||||
self.assertRaises(KeyError, inst.__exit__, KeyError, KeyError(), None)
|
||||
self.assertFalse(manager.committed)
|
||||
self.assertTrue(manager.aborted)
|
||||
|
||||
def test_explicit_mode(self):
|
||||
from .. import TransactionManager
|
||||
from ..interfaces import AlreadyInTransaction, NoTransaction
|
||||
|
||||
tm = TransactionManager()
|
||||
self.assertFalse(tm.explicit)
|
||||
|
||||
tm = TransactionManager(explicit=True)
|
||||
self.assertTrue(tm.explicit)
|
||||
for name in 'get', 'commit', 'abort', 'doom', 'isDoomed', 'savepoint':
|
||||
with self.assertRaises(NoTransaction):
|
||||
getattr(tm, name)()
|
||||
|
||||
t = tm.begin()
|
||||
with self.assertRaises(AlreadyInTransaction):
|
||||
tm.begin()
|
||||
|
||||
self.assertTrue(t is tm.get())
|
||||
|
||||
self.assertFalse(tm.isDoomed())
|
||||
tm.doom()
|
||||
self.assertTrue(tm.isDoomed())
|
||||
tm.abort()
|
||||
|
||||
for name in 'get', 'commit', 'abort', 'doom', 'isDoomed', 'savepoint':
|
||||
with self.assertRaises(NoTransaction):
|
||||
getattr(tm, name)()
|
||||
|
||||
t = tm.begin()
|
||||
self.assertFalse(tm.isDoomed())
|
||||
with self.assertRaises(AlreadyInTransaction):
|
||||
tm.begin()
|
||||
tm.savepoint()
|
||||
tm.commit()
|
||||
|
||||
|
||||
|
||||
class DummyManager(object):
|
||||
entered = False
|
||||
committed = False
|
||||
aborted = False
|
||||
|
||||
def __init__(self, raise_on_commit=None, raise_on_abort=None):
|
||||
self.raise_on_commit = raise_on_commit
|
||||
self.raise_on_abort = raise_on_abort
|
||||
|
||||
def _retryable(self, t, v):
|
||||
from transaction._manager import TransientError
|
||||
return issubclass(t, TransientError)
|
||||
|
||||
def __enter__(self):
|
||||
self.entered = True
|
||||
|
||||
def abort(self):
|
||||
self.aborted = True
|
||||
if self.raise_on_abort:
|
||||
raise self.raise_on_abort
|
||||
|
||||
def commit(self):
|
||||
self.committed = True
|
||||
if self.raise_on_commit:
|
||||
raise self.raise_on_commit
|
||||
|
||||
|
||||
class DataObject:
|
||||
|
||||
def __init__(self, transaction_manager, nost=0):
|
||||
self.transaction_manager = transaction_manager
|
||||
self.nost = nost
|
||||
self._p_jar = None
|
||||
|
||||
def modify(self, nojar=0, tracing=0):
|
||||
if not nojar:
|
||||
if self.nost:
|
||||
self._p_jar = BasicJar(tracing=tracing)
|
||||
else:
|
||||
self._p_jar = BasicJar(tracing=tracing)
|
||||
self.transaction_manager.get().join(self._p_jar)
|
||||
|
||||
|
||||
class TestTxnException(Exception):
|
||||
pass
|
||||
|
||||
|
||||
class BasicJar(object):
|
||||
|
||||
def __init__(self, errors=(), tracing=0):
|
||||
if not isinstance(errors, tuple):
|
||||
errors = errors,
|
||||
self.errors = errors
|
||||
self.tracing = tracing
|
||||
self.cabort = 0
|
||||
self.ccommit = 0
|
||||
self.ctpc_begin = 0
|
||||
self.ctpc_abort = 0
|
||||
self.ctpc_vote = 0
|
||||
self.ctpc_finish = 0
|
||||
self.cabort_sub = 0
|
||||
self.ccommit_sub = 0
|
||||
|
||||
def __repr__(self):
|
||||
return "<%s %X %s>" % (self.__class__.__name__,
|
||||
positive_id(self),
|
||||
self.errors)
|
||||
|
||||
def sortKey(self):
|
||||
# All these jars use the same sort key, and Python's list.sort()
|
||||
# is stable. These two
|
||||
return self.__class__.__name__
|
||||
|
||||
def check(self, method):
|
||||
if self.tracing: # pragma: no cover
|
||||
print('%s calling method %s'%(str(self.tracing),method))
|
||||
|
||||
if method in self.errors:
|
||||
raise TestTxnException("error %s" % method)
|
||||
|
||||
## basic jar txn interface
|
||||
|
||||
def abort(self, *args):
|
||||
self.check('abort')
|
||||
self.cabort += 1
|
||||
|
||||
def commit(self, *args):
|
||||
self.check('commit')
|
||||
self.ccommit += 1
|
||||
|
||||
def tpc_begin(self, txn, sub=0):
|
||||
self.check('tpc_begin')
|
||||
self.ctpc_begin += 1
|
||||
|
||||
def tpc_vote(self, *args):
|
||||
self.check('tpc_vote')
|
||||
self.ctpc_vote += 1
|
||||
|
||||
def tpc_abort(self, *args):
|
||||
self.check('tpc_abort')
|
||||
self.ctpc_abort += 1
|
||||
|
||||
def tpc_finish(self, *args):
|
||||
self.check('tpc_finish')
|
||||
self.ctpc_finish += 1
|
||||
|
||||
|
||||
class DummySynch(object):
|
||||
def __init__(self):
|
||||
self._txns = set()
|
||||
def newTransaction(self, txn):
|
||||
self._txns.add(txn)
|
||||
|
||||
|
||||
def positive_id(obj):
|
||||
"""Return id(obj) as a non-negative integer."""
|
||||
import struct
|
||||
_ADDRESS_MASK = 256 ** struct.calcsize('P')
|
||||
|
||||
result = id(obj)
|
||||
if result < 0: # pragma: no cover
|
||||
# Happens...on old 32-bit systems?
|
||||
result += _ADDRESS_MASK
|
||||
assert result > 0
|
||||
return result
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,142 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
##############################################################################
|
||||
"""Test backwards compatibility for resource managers using register().
|
||||
|
||||
The transaction package supports several different APIs for resource
|
||||
managers. The original ZODB3 API was implemented by ZODB.Connection.
|
||||
The Connection passed persistent objects to a Transaction's register()
|
||||
method. It's possible that third-party code also used this API, hence
|
||||
these tests that the code that adapts the old interface to the current
|
||||
API works.
|
||||
|
||||
These tests use a TestConnection object that implements the old API.
|
||||
They check that the right methods are called and in roughly the right
|
||||
order.
|
||||
"""
|
||||
import unittest
|
||||
|
||||
|
||||
class BBBTests(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
from transaction import abort
|
||||
abort()
|
||||
tearDown = setUp
|
||||
|
||||
def test_basic_commit(self):
|
||||
import transaction
|
||||
cn = TestConnection()
|
||||
cn.register(Object())
|
||||
cn.register(Object())
|
||||
cn.register(Object())
|
||||
transaction.commit()
|
||||
self.assertEqual(len(cn.committed), 3)
|
||||
self.assertEqual(len(cn.aborted), 0)
|
||||
self.assertEqual(cn.calls, ['begin', 'vote', 'finish'])
|
||||
|
||||
def test_basic_abort(self):
|
||||
# If the application calls abort(), then the transaction never gets
|
||||
# into the two-phase commit. It just aborts each object.
|
||||
import transaction
|
||||
cn = TestConnection()
|
||||
cn.register(Object())
|
||||
cn.register(Object())
|
||||
cn.register(Object())
|
||||
transaction.abort()
|
||||
self.assertEqual(len(cn.committed), 0)
|
||||
self.assertEqual(len(cn.aborted), 3)
|
||||
self.assertEqual(cn.calls, [])
|
||||
|
||||
def test_tpc_error(self):
|
||||
# The tricky part of the implementation is recovering from an error
|
||||
# that occurs during the two-phase commit. We override the commit()
|
||||
# and abort() methods of Object to cause errors during commit.
|
||||
|
||||
# Note that the implementation uses lists internally, so that objects
|
||||
# are committed in the order they are registered. (In the presence
|
||||
# of multiple resource managers, objects from a single resource
|
||||
# manager are committed in order. I'm not sure if this is an
|
||||
# accident of the implementation or a feature that should be
|
||||
# supported by any implementation.)
|
||||
|
||||
# The order of resource managers depends on sortKey().
|
||||
import transaction
|
||||
cn = TestConnection()
|
||||
cn.register(Object())
|
||||
cn.register(CommitError())
|
||||
cn.register(Object())
|
||||
self.assertRaises(RuntimeError, transaction.commit)
|
||||
self.assertEqual(len(cn.committed), 1)
|
||||
self.assertEqual(len(cn.aborted), 3)
|
||||
|
||||
|
||||
class Object(object):
|
||||
|
||||
def commit(self):
|
||||
pass
|
||||
|
||||
def abort(self):
|
||||
pass
|
||||
|
||||
|
||||
class CommitError(Object):
|
||||
|
||||
def commit(self):
|
||||
raise RuntimeError("commit")
|
||||
|
||||
|
||||
class AbortError(Object):
|
||||
|
||||
def abort(self):
|
||||
raise AssertionError("This should not actually be called")
|
||||
|
||||
|
||||
class BothError(CommitError, AbortError):
|
||||
pass
|
||||
|
||||
|
||||
class TestConnection(object):
|
||||
|
||||
def __init__(self):
|
||||
self.committed = []
|
||||
self.aborted = []
|
||||
self.calls = []
|
||||
|
||||
def register(self, obj):
|
||||
import transaction
|
||||
obj._p_jar = self
|
||||
transaction.get().register(obj)
|
||||
|
||||
def sortKey(self):
|
||||
return str(id(self))
|
||||
|
||||
def tpc_begin(self, txn):
|
||||
self.calls.append("begin")
|
||||
|
||||
def tpc_vote(self, txn):
|
||||
self.calls.append("vote")
|
||||
|
||||
def tpc_finish(self, txn):
|
||||
self.calls.append("finish")
|
||||
|
||||
def tpc_abort(self, txn):
|
||||
self.calls.append("abort")
|
||||
|
||||
def commit(self, obj, txn):
|
||||
obj.commit()
|
||||
self.committed.append(obj)
|
||||
|
||||
def abort(self, obj, txn):
|
||||
obj.abort()
|
||||
self.aborted.append(obj)
|
||||
@@ -0,0 +1,66 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2004 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
##############################################################################
|
||||
import unittest
|
||||
|
||||
|
||||
class SavepointTests(unittest.TestCase):
|
||||
|
||||
def testRollbackRollsbackDataManagersThatJoinedLater(self):
|
||||
# A savepoint needs to not just rollback it's savepoints, but needs
|
||||
# to # rollback savepoints for data managers that joined savepoints
|
||||
# after the savepoint:
|
||||
import transaction
|
||||
from transaction.tests import savepointsample
|
||||
dm = savepointsample.SampleSavepointDataManager()
|
||||
dm['name'] = 'bob'
|
||||
sp1 = transaction.savepoint()
|
||||
dm['job'] = 'geek'
|
||||
sp2 = transaction.savepoint()
|
||||
dm['salary'] = 'fun'
|
||||
dm2 = savepointsample.SampleSavepointDataManager()
|
||||
dm2['name'] = 'sally'
|
||||
|
||||
self.assertTrue('name' in dm)
|
||||
self.assertTrue('job' in dm)
|
||||
self.assertTrue('salary' in dm)
|
||||
self.assertTrue('name' in dm2)
|
||||
|
||||
sp1.rollback()
|
||||
|
||||
self.assertTrue('name' in dm)
|
||||
self.assertFalse('job' in dm)
|
||||
self.assertFalse('salary' in dm)
|
||||
self.assertFalse('name' in dm2)
|
||||
|
||||
def test_commit_after_rollback_for_dm_that_joins_after_savepoint(self):
|
||||
# There was a problem handling data managers that joined after a
|
||||
# savepoint. If the savepoint was rolled back and then changes
|
||||
# made, the dm would end up being joined twice, leading to extra
|
||||
# tpc calls and pain.
|
||||
import transaction
|
||||
from transaction.tests import savepointsample
|
||||
sp = transaction.savepoint()
|
||||
dm = savepointsample.SampleSavepointDataManager()
|
||||
dm['name'] = 'bob'
|
||||
sp.rollback()
|
||||
dm['name'] = 'Bob'
|
||||
transaction.commit()
|
||||
self.assertEqual(dm['name'], 'Bob')
|
||||
|
||||
|
||||
|
||||
def test_suite():
|
||||
return unittest.TestSuite((
|
||||
unittest.makeSuite(SavepointTests),
|
||||
))
|
||||
@@ -0,0 +1,131 @@
|
||||
##############################################################################
|
||||
#
|
||||
# Copyright (c) 2007 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE
|
||||
#
|
||||
##############################################################################
|
||||
import unittest
|
||||
from transaction._compat import JYTHON
|
||||
|
||||
class WeakSetTests(unittest.TestCase):
|
||||
def test_contains(self):
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
w.add(dummy)
|
||||
self.assertEqual(dummy in w, True)
|
||||
dummy2 = Dummy()
|
||||
self.assertEqual(dummy2 in w, False)
|
||||
|
||||
def test_len(self):
|
||||
import gc
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
d1 = Dummy()
|
||||
d2 = Dummy()
|
||||
w.add(d1)
|
||||
w.add(d2)
|
||||
self.assertEqual(len(w), 2)
|
||||
del d1
|
||||
gc.collect()
|
||||
if not JYTHON:
|
||||
# The Jython GC is non deterministic
|
||||
self.assertEqual(len(w), 1)
|
||||
|
||||
def test_remove(self):
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
w.add(dummy)
|
||||
self.assertEqual(dummy in w, True)
|
||||
w.remove(dummy)
|
||||
self.assertEqual(dummy in w, False)
|
||||
|
||||
def test_clear(self):
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
w.add(dummy)
|
||||
dummy2 = Dummy()
|
||||
w.add(dummy2)
|
||||
self.assertEqual(dummy in w, True)
|
||||
self.assertEqual(dummy2 in w, True)
|
||||
w.clear()
|
||||
self.assertEqual(dummy in w, False)
|
||||
self.assertEqual(dummy2 in w, False)
|
||||
|
||||
def test_as_weakref_list(self):
|
||||
import gc
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
dummy2 = Dummy()
|
||||
dummy3 = Dummy()
|
||||
w.add(dummy)
|
||||
w.add(dummy2)
|
||||
w.add(dummy3)
|
||||
del dummy3
|
||||
gc.collect()
|
||||
refs = w.as_weakref_list()
|
||||
self.assertTrue(isinstance(refs, list))
|
||||
L = [x() for x in refs]
|
||||
# L is a list, but it does not have a guaranteed order.
|
||||
self.assertTrue(list, type(L))
|
||||
self.assertEqual(set(L), set([dummy, dummy2]))
|
||||
|
||||
def test_map(self):
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
dummy2 = Dummy()
|
||||
dummy3 = Dummy()
|
||||
w.add(dummy)
|
||||
w.add(dummy2)
|
||||
w.add(dummy3)
|
||||
def poker(x):
|
||||
x.poked = 1
|
||||
w.map(poker)
|
||||
for thing in dummy, dummy2, dummy3:
|
||||
self.assertEqual(thing.poked, 1)
|
||||
|
||||
def test_map_w_gced_element(self):
|
||||
import gc
|
||||
from transaction.weakset import WeakSet
|
||||
w = WeakSet()
|
||||
dummy = Dummy()
|
||||
dummy2 = Dummy()
|
||||
dummy3 = [Dummy()]
|
||||
w.add(dummy)
|
||||
w.add(dummy2)
|
||||
w.add(dummy3[0])
|
||||
|
||||
_orig = w.as_weakref_list
|
||||
def _as_weakref_list():
|
||||
# simulate race condition during iteration of list
|
||||
# object is collected after being iterated.
|
||||
result = _orig()
|
||||
del dummy3[:]
|
||||
gc.collect()
|
||||
return result
|
||||
w.as_weakref_list = _as_weakref_list
|
||||
|
||||
def poker(x):
|
||||
x.poked = 1
|
||||
w.map(poker)
|
||||
for thing in dummy, dummy2:
|
||||
self.assertEqual(thing.poked, 1)
|
||||
|
||||
|
||||
class Dummy:
|
||||
pass
|
||||
|
||||
|
||||
def test_suite():
|
||||
return unittest.makeSuite(WeakSetTests)
|
||||
@@ -0,0 +1,85 @@
|
||||
############################################################################
|
||||
#
|
||||
# Copyright (c) 2007 Zope Foundation and Contributors.
|
||||
# All Rights Reserved.
|
||||
#
|
||||
# This software is subject to the provisions of the Zope Public License,
|
||||
# Version 2.1 (ZPL). A copy of the ZPL should accompany this distribution.
|
||||
# THIS SOFTWARE IS PROVIDED "AS IS" AND ANY AND ALL EXPRESS OR IMPLIED
|
||||
# WARRANTIES ARE DISCLAIMED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||
# WARRANTIES OF TITLE, MERCHANTABILITY, AGAINST INFRINGEMENT, AND FITNESS
|
||||
# FOR A PARTICULAR PURPOSE.
|
||||
#
|
||||
############################################################################
|
||||
|
||||
import weakref
|
||||
|
||||
# A simple implementation of weak sets, supplying just enough of Python's
|
||||
# sets.Set interface for our needs.
|
||||
|
||||
class WeakSet(object):
|
||||
"""A set of objects that doesn't keep its elements alive.
|
||||
|
||||
The objects in the set must be weakly referencable.
|
||||
The objects need not be hashable, and need not support comparison.
|
||||
Two objects are considered to be the same iff their id()s are equal.
|
||||
|
||||
When the only references to an object are weak references (including
|
||||
those from WeakSets), the object can be garbage-collected, and
|
||||
will vanish from any WeakSets it may be a member of at that time.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
# Map id(obj) to obj. By using ids as keys, we avoid requiring
|
||||
# that the elements be hashable or comparable.
|
||||
self.data = weakref.WeakValueDictionary()
|
||||
|
||||
def __len__(self):
|
||||
return len(self.data)
|
||||
|
||||
def __contains__(self, obj):
|
||||
return id(obj) in self.data
|
||||
|
||||
# Same as a Set, add obj to the collection.
|
||||
def add(self, obj):
|
||||
self.data[id(obj)] = obj
|
||||
|
||||
# Same as a Set, remove obj from the collection, and raise
|
||||
# KeyError if obj not in the collection.
|
||||
def remove(self, obj):
|
||||
del self.data[id(obj)]
|
||||
|
||||
def clear(self):
|
||||
self.data.clear()
|
||||
|
||||
# f is a one-argument function. Execute f(elt) for each elt in the
|
||||
# set. f's return value is ignored.
|
||||
def map(self, f):
|
||||
for wr in self.as_weakref_list():
|
||||
elt = wr()
|
||||
if elt is not None:
|
||||
f(elt)
|
||||
|
||||
# Return a list of weakrefs to all the objects in the collection.
|
||||
# Because a weak dict is used internally, iteration is dicey (the
|
||||
# underlying dict may change size during iteration, due to gc or
|
||||
# activity from other threads). as_weakef_list() is safe.
|
||||
#
|
||||
# If we invoke self.data.values() instead, we get back a list of live
|
||||
# objects instead of weakrefs. If gc occurs while this list is alive,
|
||||
# all the objects move to an older generation (because they're strongly
|
||||
# referenced by the list!). They can't get collected then, until a
|
||||
# less frequent collection of the older generation. Before then, if we
|
||||
# invoke self.data.values() again, they're still alive, and if gc occurs
|
||||
# while that list is alive they're all moved to yet an older generation.
|
||||
# And so on. Stress tests showed that it was easy to get into a state
|
||||
# where a WeakSet grows without bounds, despite that almost all its
|
||||
# elements are actually trash. By returning a list of weakrefs instead,
|
||||
# we avoid that, although the decision to use weakrefs is now very
|
||||
# visible to our clients.
|
||||
|
||||
def as_weakref_list(self):
|
||||
# The docstring of WeakValueDictionary.valuerefs()
|
||||
# guarantees to return an actual list on all supported versions
|
||||
# of Python.
|
||||
return self.data.valuerefs()
|
||||
Reference in New Issue
Block a user