forked from weaselab/conflict-set
Compare commits
16
Commits
f30887f280
...
dee3a8f640
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dee3a8f640 | ||
|
|
732d19efa1 | ||
|
|
9449190d02 | ||
|
|
4fcdc5d7e9 | ||
|
|
63f9a139da | ||
|
|
8f9f345c64 | ||
|
|
b9b2d69dd5 | ||
|
|
d3c8f4afc6 | ||
|
|
549724f09e | ||
|
|
52eb13cc0b | ||
|
|
ccd637deab | ||
|
|
971deb477c | ||
|
|
ff0722728a | ||
|
|
60881419b8 | ||
|
|
4515af3662 | ||
|
|
7eaac2a184 |
@@ -37,15 +37,16 @@ ConflictSet::ReadRange singleton(Arena &arena, TrivialSpan key) {
|
||||
}
|
||||
|
||||
ConflictSet::ReadRange prefixRange(Arena &arena, TrivialSpan key) {
|
||||
int index;
|
||||
for (index = key.size() - 1; index >= 0; index--)
|
||||
if ((key[index]) != 255)
|
||||
int index = key.size() - 1;
|
||||
for (; index >= 0; index--)
|
||||
if (key[index] != 255)
|
||||
break;
|
||||
|
||||
// Must not be called with a string that consists only of zero or more '\xff'
|
||||
// bytes.
|
||||
// bytes, or with an empty string (which has no finite upper bound).
|
||||
if (index < 0) {
|
||||
assert(false);
|
||||
std::abort();
|
||||
}
|
||||
|
||||
uint8_t *buf = new (arena) uint8_t[index + 1];
|
||||
|
||||
@@ -139,6 +139,7 @@ add_custom_command(
|
||||
COMMAND_EXPAND_LISTS)
|
||||
|
||||
add_library(${PROJECT_NAME} SHARED ${CMAKE_BINARY_DIR}/${PROJECT_NAME}.o)
|
||||
add_dependencies(${PROJECT_NAME} ${PROJECT_NAME}-object)
|
||||
set_target_properties(
|
||||
${PROJECT_NAME} PROPERTIES LIBRARY_OUTPUT_DIRECTORY
|
||||
"${CMAKE_CURRENT_BINARY_DIR}/radix_tree")
|
||||
@@ -155,6 +156,7 @@ if(HAS_VERSION_SCRIPT)
|
||||
endif()
|
||||
|
||||
add_library(${PROJECT_NAME}-static STATIC ${CMAKE_BINARY_DIR}/${PROJECT_NAME}.o)
|
||||
add_dependencies(${PROJECT_NAME}-static ${PROJECT_NAME}-object)
|
||||
if(CMAKE_BUILD_TYPE STREQUAL Debug)
|
||||
set_target_properties(${PROJECT_NAME}-static PROPERTIES LINKER_LANGUAGE CXX)
|
||||
else()
|
||||
|
||||
+6
-5
@@ -5679,13 +5679,13 @@ std::string getPartialKeyPrintable(Node *n) {
|
||||
}
|
||||
|
||||
std::string strinc(std::string_view str, bool &ok) {
|
||||
int index;
|
||||
for (index = str.size() - 1; index >= 0; index--)
|
||||
if ((uint8_t &)(str[index]) != 255)
|
||||
int index = static_cast<int>(str.size()) - 1;
|
||||
for (; index >= 0; index--)
|
||||
if (static_cast<uint8_t>(str[index]) != 255)
|
||||
break;
|
||||
|
||||
// Must not be called with a string that consists only of zero or more
|
||||
// '\xff' bytes.
|
||||
// '\xff' bytes, and the empty string has no successor.
|
||||
if (index < 0) {
|
||||
ok = false;
|
||||
return {};
|
||||
@@ -5693,7 +5693,8 @@ std::string strinc(std::string_view str, bool &ok) {
|
||||
ok = true;
|
||||
|
||||
auto r = std::string(str.substr(0, index + 1));
|
||||
((uint8_t &)r[r.size() - 1])++;
|
||||
auto &last = r[r.size() - 1];
|
||||
last = static_cast<char>(static_cast<uint8_t>(last) + 1);
|
||||
return r;
|
||||
}
|
||||
|
||||
|
||||
+8
-3
@@ -96,7 +96,9 @@ void ConflictSet::setOldestVersion(int64_t oldestVersion) {
|
||||
return impl->setOldestVersion(oldestVersion);
|
||||
}
|
||||
|
||||
int64_t ConflictSet::getBytes() const { return -1; }
|
||||
// The hash_table implementation does not track memory usage, so return 0 to
|
||||
// satisfy the API contract that getBytes() returns a non-negative value.
|
||||
int64_t ConflictSet::getBytes() const { return 0; }
|
||||
|
||||
void ConflictSet::getMetricsV1(MetricsV1 **metrics, int *count) const {
|
||||
*metrics = nullptr;
|
||||
@@ -161,7 +163,10 @@ __attribute__((__visibility__("default"))) void ConflictSet_destroy(void *cs) {
|
||||
}
|
||||
__attribute__((__visibility__("default"))) int64_t
|
||||
ConflictSet_getBytes(void *cs) {
|
||||
using Impl = ConflictSet::Impl;
|
||||
return -1;
|
||||
(void)cs;
|
||||
// The hash_table implementation does not track memory usage, so return 0 to
|
||||
// satisfy the API contract that ConflictSet_getBytes returns a non-negative
|
||||
// value.
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
+7
-6
@@ -1,5 +1,6 @@
|
||||
#include <ConflictSet.h>
|
||||
|
||||
#include <algorithm>
|
||||
#include <cerrno>
|
||||
#include <chrono>
|
||||
#include <cstdio>
|
||||
@@ -77,10 +78,10 @@ int main(int argc, const char **argv) {
|
||||
begin = end + 1;
|
||||
end = (uint8_t *)memchr(begin, '\n', size);
|
||||
|
||||
if (line.size() > 0 && line[0] == 'P') {
|
||||
write = line.subspan(2, line.size());
|
||||
} else if (line.size() > 0 && line[0] == 'L') {
|
||||
reads.push_back(line.subspan(2, line.size()));
|
||||
if (line.size() >= 2 && line[0] == 'P') {
|
||||
write = line.subspan(2, line.size() - 2);
|
||||
} else if (line.size() >= 2 && line[0] == 'L') {
|
||||
reads.push_back(line.subspan(2, line.size() - 2));
|
||||
} else if (line.empty()) {
|
||||
{
|
||||
readRanges.resize(reads.size());
|
||||
@@ -90,7 +91,7 @@ int main(int argc, const char **argv) {
|
||||
iter->begin.len = read.size();
|
||||
checkBytes += read.size();
|
||||
iter->end.len = 0;
|
||||
iter->readVersion = version - 100;
|
||||
iter->readVersion = std::max<int64_t>(0, version - 100);
|
||||
++iter;
|
||||
}
|
||||
}
|
||||
@@ -121,7 +122,7 @@ int main(int argc, const char **argv) {
|
||||
}
|
||||
|
||||
timer = now();
|
||||
cs.setOldestVersion(version - 10000);
|
||||
cs.setOldestVersion(std::max<int64_t>(0, version - 10000));
|
||||
gcTime += now() - timer;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
#include <sys/ioctl.h>
|
||||
#include <sys/resource.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/syscall.h>
|
||||
#include <sys/types.h>
|
||||
#include <sys/uio.h>
|
||||
#include <thread>
|
||||
|
||||
+23
-9
@@ -27,23 +27,37 @@ class Result(enum.Enum):
|
||||
TOO_OLD = 2
|
||||
|
||||
|
||||
def write(begin: bytes, end: Optional[bytes] = None) -> WriteRange:
|
||||
b = (ctypes.c_ubyte * len(begin)).from_buffer(bytearray(begin))
|
||||
def _make_key(buf: bytes) -> tuple[_Key, bytearray]:
|
||||
"""Create a _Key and a backing bytearray that must be kept alive."""
|
||||
backing = bytearray(buf)
|
||||
array = (ctypes.c_ubyte * len(backing)).from_buffer(backing)
|
||||
return _Key(array, len(array)), backing
|
||||
|
||||
|
||||
def write(begin: bytes, end: Optional[bytes] = None) -> WriteRange:
|
||||
begin_key, begin_buf = _make_key(begin)
|
||||
if end is None:
|
||||
e = (ctypes.c_ubyte * 0)()
|
||||
end_key = _Key((ctypes.c_ubyte * 0)(), 0)
|
||||
end_buf = None
|
||||
else:
|
||||
e = (ctypes.c_ubyte * len(end)).from_buffer(bytearray(end))
|
||||
return WriteRange(_Key(b, len(b)), _Key(e, len(e)))
|
||||
end_key, end_buf = _make_key(end)
|
||||
result = WriteRange(begin_key, end_key)
|
||||
result._begin_buf = begin_buf
|
||||
result._end_buf = end_buf
|
||||
return result
|
||||
|
||||
|
||||
def read(version: int, begin: bytes, end: Optional[bytes] = None) -> ReadRange:
|
||||
b = (ctypes.c_ubyte * len(begin)).from_buffer(bytearray(begin))
|
||||
begin_key, begin_buf = _make_key(begin)
|
||||
if end is None:
|
||||
e = (ctypes.c_ubyte * 0)()
|
||||
end_key = _Key((ctypes.c_ubyte * 0)(), 0)
|
||||
end_buf = None
|
||||
else:
|
||||
e = (ctypes.c_ubyte * len(end)).from_buffer(bytearray(end))
|
||||
return ReadRange(_Key(b, len(b)), _Key(e, len(e)), version)
|
||||
end_key, end_buf = _make_key(end)
|
||||
result = ReadRange(begin_key, end_key, version)
|
||||
result._begin_buf = begin_buf
|
||||
result._end_buf = end_buf
|
||||
return result
|
||||
|
||||
|
||||
class ConflictSet:
|
||||
|
||||
+11
-2
@@ -88,7 +88,8 @@ struct __attribute__((__visibility__("default"))) ConflictSet {
|
||||
|
||||
~ConflictSet();
|
||||
|
||||
/** Returns the total bytes in use by this ConflictSet */
|
||||
/** Returns the total bytes in use by this ConflictSet. Implementations that
|
||||
* do not track memory usage return 0. */
|
||||
int64_t getBytes() const;
|
||||
|
||||
/** Experimental! */
|
||||
@@ -132,6 +133,13 @@ struct __attribute__((__visibility__("default"))) ConflictSet {
|
||||
|
||||
private:
|
||||
Impl *impl;
|
||||
#if __cplusplus <= 199711L
|
||||
/* Declared private and left undefined to prevent copying in C++98/C++03.
|
||||
The compiler would otherwise implicitly generate public copy operations,
|
||||
which share the opaque Impl* and cause a double-free. */
|
||||
ConflictSet(const ConflictSet &);
|
||||
ConflictSet &operator=(const ConflictSet &);
|
||||
#endif
|
||||
};
|
||||
} /* namespace weaselab */
|
||||
|
||||
@@ -211,7 +219,8 @@ ConflictSet *ConflictSet_create(int64_t oldestVersion);
|
||||
|
||||
void ConflictSet_destroy(ConflictSet *cs);
|
||||
|
||||
/** Returns the total bytes in use by this ConflictSet */
|
||||
/** Returns the total bytes in use by this ConflictSet. Implementations that
|
||||
* do not track memory usage return 0. */
|
||||
int64_t ConflictSet_getBytes(const ConflictSet *cs);
|
||||
|
||||
#endif
|
||||
|
||||
+49
-2
@@ -57,6 +57,53 @@ def test_conflict_set():
|
||||
assert cs.check(read(0, key), read(1, key)) == [Result.TOO_OLD, Result.COMMIT]
|
||||
|
||||
|
||||
def test_hash_table_getBytes():
|
||||
# Regression test for issue #62: the hash_table implementation is
|
||||
# point-query only and does not track memory usage, but getBytes() must
|
||||
# still return a non-negative value rather than -1.
|
||||
with ConflictSet(0, build_dir=build_dir, implementation="hash_table") as cs:
|
||||
assert cs.getBytes() == 0
|
||||
cs.addWrites(1, write(b"key"))
|
||||
assert cs.getBytes() >= 0
|
||||
assert cs.check(read(0, b"key")) == [Result.CONFLICT]
|
||||
|
||||
|
||||
def test_write_read_without_outer_reference():
|
||||
# Regression test for issue #42: WriteRange/ReadRange must keep their
|
||||
# backing key buffers alive, because the C library reads the pointer
|
||||
# stored in _Key while addWrites/check run.
|
||||
with DebugConflictSet() as cs:
|
||||
# The bytes literal is not referenced after this expression.
|
||||
cs.addWrites(1, write(b"key"))
|
||||
assert cs.check(read(0, b"key")) == [Result.CONFLICT]
|
||||
|
||||
cs.addWrites(2, write(b"a", b"z"))
|
||||
assert cs.check(read(1, b"a", b"z")) == [Result.CONFLICT]
|
||||
assert cs.check(read(1, b"b")) == [Result.CONFLICT]
|
||||
assert cs.check(read(1, b"0")) == [Result.COMMIT]
|
||||
|
||||
|
||||
def test_range_keeps_key_buffers_alive():
|
||||
# Verify the fix for issue #42: returned range objects must retain a
|
||||
# reference to the backing bytearray so the C pointer stays valid after
|
||||
# the helper returns.
|
||||
w = write(b"key")
|
||||
assert w._begin_buf == bytearray(b"key")
|
||||
assert w._end_buf is None
|
||||
|
||||
w2 = write(b"a", b"z")
|
||||
assert w2._begin_buf == bytearray(b"a")
|
||||
assert w2._end_buf == bytearray(b"z")
|
||||
|
||||
r = read(0, b"key")
|
||||
assert r._begin_buf == bytearray(b"key")
|
||||
assert r._end_buf is None
|
||||
|
||||
r2 = read(1, b"a", b"z")
|
||||
assert r2._begin_buf == bytearray(b"a")
|
||||
assert r2._end_buf == bytearray(b"z")
|
||||
|
||||
|
||||
def test_update_zero_should_commit():
|
||||
with DebugConflictSet() as cs1:
|
||||
with DebugConflictSet() as cs2:
|
||||
@@ -68,7 +115,7 @@ def test_update_zero_should_commit():
|
||||
for i in range(256 - 17, 256):
|
||||
cs2.addWrites(int(1), write(bytes([i])))
|
||||
# Scan until first point write
|
||||
cs2.check(read(0, b"\x00", bytes([256 - 17])))
|
||||
assert cs2.check(read(0, b"\x00", bytes([256 - 17]))) == [Result.COMMIT]
|
||||
|
||||
|
||||
def test_update_zero_should_conflict():
|
||||
@@ -81,7 +128,7 @@ def test_update_zero_should_conflict():
|
||||
# "zero" is now 2**31 + 100
|
||||
cs1.addWrites(2**32 + 101, write(b"", b"\x02"), write(b"\x01"))
|
||||
# rangeVersion of \x01 is now 2**31 + 100 ("max" of (2**31 + 100, 2**32 + 101))
|
||||
cs1.check(read(2**32 + 1, b"\x00"))
|
||||
assert cs1.check(read(2**32 + 1, b"\x00")) == [Result.CONFLICT]
|
||||
# but 2**32 + 1 ">" 2**31 + 100 , and it incorrectly commits
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user