openCursor0(ReadableTransaction txn, boolean partOfAWholeTreeWalk)
{
return transformKeysAndValues(
partOfAWholeTreeWalk ? txn.openBulkCursor(getName()) : txn.openCursor(getName()), TO_KEY, TO_LONG);
}
void addCount(final WriteableTransaction txn, ByteSequence key, final long delta)
{
txn.update(getName(), getShardedKey(key), new UpdateFunction()
{
@Override
public ByteSequence computeNewValue(ByteSequence oldValue)
{
final long currentValue = oldValue == null ? 0 : decodeValue(oldValue.toByteString());
return encodeValue(currentValue + delta);
}
});
}
void importPut(Importer importer, ByteSequence key, long delta)
{
if (delta != 0)
{
importer.put(getName(), getShardedKey(key), encodeValue(delta));
}
}
long getCount(final ReadableTransaction txn, ByteSequence key)
{
return getCount(txn, key, false);
}
/**
* The same read, told which kind of work it is part of. A client operation reads a counter of its
* own and takes the bound of one - {@code numSubordinates} of a search
* ({@code EntryContainer.getNumberOfChildren}), the entry count of a VLV index a search is paging
* through ({@code VLVIndex.getEntryCount}) - while {@code verify-index} reads one per DN of the
* tree it is walking, with nobody waiting on it: bounding those as client operations is what #877
* exists to stop, and on the JDBC backend it aborted a verify of a backend large enough.
*
* The third caller is {@code ID2ChildrenCount.getTotalCount}, which is read both ways and is told
* which it is by its own caller: a verify sizes its progress report with it, {@code cn=monitor}
* and the searches of {@code GroupManager} and {@code SubentryManager} read it for a client. A
* delete and a modify DN reach neither form - they go through {@link #removeCount}, which is a
* client operation by construction.
*
* @param txn storage transaction
* @param key the counter to read
* @param partOfAWholeTreeWalk whether this read belongs to a walk of a whole tree rather than to
* a client operation
* @return Value of the counter. 0 if no counter is associated yet.
* @see ReadableTransaction#openBulkCursor(TreeName)
*/
long getCount(final ReadableTransaction txn, ByteSequence key, boolean partOfAWholeTreeWalk)
{
long counterValue = 0;
try (final SequentialCursor cursor =
new ShardCursor(openCursor0(txn, partOfAWholeTreeWalk), key))
{
while (cursor.next())
{
counterValue += cursor.getValue();
}
}
return counterValue;
}
long removeCount(final WriteableTransaction txn, ByteSequence key)
{
long counterValue = 0;
// a removal is always a client operation: an entry is being deleted or moved
try (final SequentialCursor cursor = new ShardCursor(openCursor0(txn, false), key))
{
// Iterate over and remove all the thread local shards
while (cursor.next())
{
counterValue += cursor.getValue();
cursor.delete();
}
}
return counterValue;
}
static long decodeValue(ByteString value)
{
switch (value.length())
{
case 1:
return value.byteAt(0);
case (Integer.SIZE / Byte.SIZE):
return value.toInt();
case (Long.SIZE / Byte.SIZE):
return value.toLong();
default:
throw new IllegalArgumentException("Unsupported sharded-counter value format.");
}
}
static ByteString encodeValue(long value)
{
final byte valueAsByte = (byte) value;
if (valueAsByte == value)
{
return ByteString.wrap(new byte[] { valueAsByte });
}
final int valueAsInt = (int) value;
if (valueAsInt == value)
{
return ByteString.valueOfInt(valueAsInt);
}
return ByteString.valueOfLong(value);
}
@Override
public String valueToString(ByteString value)
{
return String.valueOf(decodeValue(value));
}
private static ByteSequence getShardedKey(ByteSequence key)
{
final byte bucket = (byte) (Thread.currentThread().getId() & (SHARD_COUNT - 1));
return new ByteStringBuilder(key.length() + ByteStringBuilder.MAX_COMPACT_SIZE).appendBytes(key).appendByte(bucket);
}
/** Restricts a cursor to the shards of a specific key. */
private final class ShardCursor extends SequentialCursorDecorator, ByteString, Long>
{
private final ByteSequence targetKey;
private boolean initialized;
ShardCursor(Cursor delegate, ByteSequence targetKey)
{
super(delegate);
this.targetKey = targetKey;
}
@Override
public boolean next()
{
if (!initialized)
{
initialized = true;
return delegate.positionToKeyOrNext(targetKey) && isOnTargetKey();
}
return delegate.next() && isOnTargetKey();
}
private boolean isOnTargetKey()
{
return targetKey.equals(delegate.getKey());
}
}
/**
* Cursor that returns unique keys and null values. Ensure that {@link #getKey()} will return a different key after
* each {@link #next()}.
*/
private static final class UniqueKeysCursor implements SequentialCursor
{
private final Cursor delegate;
private boolean isDefined;
private K key;
private UniqueKeysCursor(Cursor cursor)
{
this.delegate = cursor;
if (!delegate.isDefined())
{
delegate.next();
}
}
@Override
public boolean next()
{
isDefined = delegate.isDefined();
if (isDefined)
{
key = delegate.getKey();
skipEntriesWithSameKey();
}
return isDefined;
}
private void skipEntriesWithSameKey()
{
throwIfUndefined(this);
while (delegate.next() && key.equals(delegate.getKey()))
{
// Skip all entries having the same key.
}
// Delegate is one step beyond. When delegate.isDefined() return false, we have to return true once more.
isDefined = true;
}
@Override
public boolean isDefined()
{
return isDefined;
}
@Override
public K getKey() throws NoSuchElementException
{
throwIfUndefined(this);
return key;
}
@Override
public Void getValue() throws NoSuchElementException
{
throwIfUndefined(this);
return null;
}
@Override
public void delete() throws NoSuchElementException, UnsupportedOperationException
{
throw new UnsupportedOperationException();
}
@Override
public void close()
{
key = null;
delegate.close();
}
private static void throwIfUndefined(SequentialCursor, ?> cursor)
{
if (!cursor.isDefined())
{
throw new NoSuchElementException();
}
}
}
}