From b16520b75fb922e6ad7cfe17cf3a06424f49b70c Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Wed, 18 Mar 2026 22:09:53 +0000 Subject: [PATCH 1/7] add UCX setup file Signed-off-by: bob-021206 --- SETUP_UCX_TCP.md | 219 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 219 insertions(+) create mode 100644 SETUP_UCX_TCP.md diff --git a/SETUP_UCX_TCP.md b/SETUP_UCX_TCP.md new file mode 100644 index 0000000..e977a5d --- /dev/null +++ b/SETUP_UCX_TCP.md @@ -0,0 +1,219 @@ +# PrisKV: From Zero to Working (UCX TCP + Python client) + +This guide targets developer environments without RDMA hardware and uses UCX over TCP for connectivity. +The goal is: +1. Verify `import priskv` works (Python client is installed) +2. Connect using `PriskvClient` to a running `priskv-server` +3. Optionally validate end-to-end `set/get` + +> Note: PrisKV may require UCX-version-specific compatibility patches. If your system uses UCX 1.12, verify the items in **Section 7** exist in your source tree. + +--- + +## 0. Prerequisites (Debian/Ubuntu) + +### System dependencies + +```bash +apt-get update +apt-get install -y \ + git gcc make cmake \ + librdmacm-dev rdma-core libibverbs-dev \ + libncurses5-dev libmount-dev libevent-dev libssl-dev \ + dpkg-dev debhelper \ + pkg-config \ + python3-pybind11 python3-dev python3-pip \ + libonig-dev libhiredis-dev liburing-dev +``` + +### Python build tooling + +```bash +pip3 install pybind11 yapf==0.32.0 +``` + +If you are running as `root` inside a container, `sudo` is usually not required. + +--- + +## 1. Prepare the source tree + +```bash +cd /workspace/priskv +``` + +If the repository uses submodules (e.g., `json-c`), initialize them: + +```bash +git submodule update --init --recursive +``` + +--- + +## 2. Build `priskv-server` (C/C++) + +Recommended: build everything once, then rebuild server if needed: + +```bash +cd /workspace/priskv/PrisKV +make +``` + +Or build only the server: + +```bash +cd /workspace/priskv/PrisKV +make server +``` + +--- + +## 3. Build and install the Python client (`pypriskv`) + +If you use a virtual environment (recommended), activate it first: + +```bash +source /workspace/priskv/.venv/bin/activate +pip install -U pip setuptools wheel +``` + +### Option A: Editable install (recommended for development) + +```bash +cd /workspace/priskv/PrisKV +make all +cd pypriskv +pip install -v -e . +``` + +### Option B: Install from wheel + +```bash +cd /workspace/priskv/PrisKV +make all +cd pypriskv +python3 setup.py build_ext bdist_wheel +pip install ./dist/*.whl +``` + +--- + +## 4. Configure runtime environment variables for UCX TCP + +Make sure both server and client use the same transport configuration: + +- `PRISKV_TRANSPORT=ucx` +- `UCX_TLS=tcp` +- If using direct mode (no Redis meta service): `PRISKV_CLIENT_DIRECT_MODE=y` + +Example: + +```bash +export PRISKV_TRANSPORT=ucx +export UCX_TLS=tcp +export PRISKV_CLIENT_DIRECT_MODE=y +``` + +Optional debugging: + +```bash +export PRISKV_LOG_LEVEL=notice +``` + +--- + +## 5. Start `priskv-server` (UCX TCP) + +```bash +cd /workspace/priskv/PrisKV +export PRISKV_TRANSPORT=ucx +export UCX_TLS=tcp +export PRISKV_CLIENT_DIRECT_MODE=y +export PRISKV_USE_SHM=n + +./server/priskv-server -a 127.0.0.1 -p 6379 --acl any +``` + +Expected server log includes: +- `UCX: <...> ready` + +--- + +## 6. Verify connectivity from Python + +Activate your venv: + +```bash +source /workspace/priskv/.venv/bin/activate +``` + +Then run a minimal connectivity check: + +```bash +export PRISKV_TRANSPORT=ucx +export UCX_TLS=tcp +export PRISKV_CLIENT_DIRECT_MODE=y + +python - <<'PY' +import priskv + +c = priskv.PriskvClient("127.0.0.1", 6379, "kvcache-redis") +print("connected") +c.close() +print("closed") +PY +``` + +Optional: validate `set/get`: + +```bash +python - <<'PY' +import priskv + +c = priskv.PriskvClient("127.0.0.1", 6379, "kvcache-redis") +print("setstr:", c.setstr("k1", "v1", 2000)) +print("getstr:", c.getstr("k1")) +c.close() +PY +``` + +--- + +## 7. UCX 1.12 compatibility patches to verify + +If your system UCX is 1.12, PrisKV may need compatibility adjustments to avoid build errors or protocol/ABI mismatches. +Verify the following items exist in your checkout. + +### 7.1 `PrisKV/lib/config.c`: UCX config parser API compatibility + +- Add `` (avoid implicit `strcmp`) +- Fix `ucs_config_parser_print_opts` call signature (argument count) +- Fix wrapper signatures/arguments for `ucs_config_parser_fill_opts` and `ucs_config_parser_set_value` + +### 7.2 `PrisKV/include/priskv-config.h`: config table struct compatibility + +- In `PRISKV_CONFIG_DECLARE_TABLE`, remove `.flags = 0` if the field does not exist in UCX 1.12. + +### 7.3 `PrisKV/lib/ucx.c`: UCX 1.12 API compatibility + +- Replace packed RKEY release logic with UCX 1.12 behavior (`ucp_rkey_buffer_release()`) +- Adjust `ucp_worker_get_address()` length argument to `size_t*` and convert back to PrisKV types as needed + +### 7.4 `PrisKV/lib/ucx.c` (client/server UCX wrapper): tag-recv completion info + +- In `priskv_ucx_post_tag_recv()`, ensure `ucp_request_param_t.op_attr_mask` includes `UCP_OP_ATTR_FIELD_RECV_INFO` + +If you still see messages like `UCX: recv <...>, expected 48` or endpoint timeouts, continue deeper protocol-level debugging (request/response completion info and response struct decoding). + +--- + +## 8. Readiness checklist for developers + +- Successful `import priskv` indicates the Python extension is built/installed correctly. +- Server logs show `UCX ... ready` indicates UCX transport is initialized. +- Client logs show `established` indicates handshake/connection succeeded. +- If `set/get` hangs, investigate the end-to-end request/response path: + - server response generation + - client response receive callback + - protocol struct decoding and completion-info handling + From c04ea57d76c82dde47139339403c33ef02b8cc60 Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Wed, 18 Mar 2026 23:10:30 +0000 Subject: [PATCH 2/7] UCX TCP debug --- SETUP_UCX_TCP.md | 20 ++++---- client/Makefile | 9 +++- client/benchmark.c | 5 +- client/client.c | 6 ++- client/priskv.h | 12 ++--- client/sync.c | 7 ++- client/transport/transport.c | 15 +++--- client/transport/ucx.c | 31 ++++++------ cluster/client/Makefile | 2 +- cluster/client/client.c | 40 +++++++-------- cluster/client/client.h | 9 ++-- cluster/client/test_status.c | 71 +++++++++++++++++++++++++++ include/priskv-utils.h | 49 +++++++++++++++++++ lib/config.c | 8 +-- lib/ucx.c | 12 +++-- pypriskv/priskv/priskv_client.py | 7 +-- pypriskv/pybind.cpp | 22 ++++----- pypriskv/testing.py | 9 ++-- run_e2e_test.py | 6 ++- run_unit_test.py | 2 +- server/Makefile | 8 ++- server/kv.c | 83 ++++++++++++++++++++++++++++++-- server/kv.h | 3 +- server/test/Makefile | 4 +- server/transport/transport.c | 55 +++++++++++---------- server/transport/ucx.c | 34 +++++-------- 26 files changed, 370 insertions(+), 159 deletions(-) create mode 100644 cluster/client/test_status.c diff --git a/SETUP_UCX_TCP.md b/SETUP_UCX_TCP.md index e977a5d..0404074 100644 --- a/SETUP_UCX_TCP.md +++ b/SETUP_UCX_TCP.md @@ -29,6 +29,8 @@ apt-get install -y \ ### Python build tooling ```bash +python3 -m venv .venv +source .venv/bin/activate pip3 install pybind11 yapf==0.32.0 ``` @@ -39,7 +41,7 @@ If you are running as `root` inside a container, `sudo` is usually not required. ## 1. Prepare the source tree ```bash -cd /workspace/priskv +cd /PrisKV ``` If the repository uses submodules (e.g., `json-c`), initialize them: @@ -55,14 +57,14 @@ git submodule update --init --recursive Recommended: build everything once, then rebuild server if needed: ```bash -cd /workspace/priskv/PrisKV +cd /PrisKV make ``` -Or build only the server: +Or build only the server (preferred if build without RDMA): ```bash -cd /workspace/priskv/PrisKV +cd /PrisKV make server ``` @@ -73,14 +75,14 @@ make server If you use a virtual environment (recommended), activate it first: ```bash -source /workspace/priskv/.venv/bin/activate +source .venv/bin/activate pip install -U pip setuptools wheel ``` ### Option A: Editable install (recommended for development) ```bash -cd /workspace/priskv/PrisKV +cd /PrisKV make all cd pypriskv pip install -v -e . @@ -89,7 +91,7 @@ pip install -v -e . ### Option B: Install from wheel ```bash -cd /workspace/priskv/PrisKV +cd /PrisKV make all cd pypriskv python3 setup.py build_ext bdist_wheel @@ -125,7 +127,7 @@ export PRISKV_LOG_LEVEL=notice ## 5. Start `priskv-server` (UCX TCP) ```bash -cd /workspace/priskv/PrisKV +cd /PrisKV export PRISKV_TRANSPORT=ucx export UCX_TLS=tcp export PRISKV_CLIENT_DIRECT_MODE=y @@ -144,7 +146,7 @@ Expected server log includes: Activate your venv: ```bash -source /workspace/priskv/.venv/bin/activate +source .venv/bin/activate ``` Then run a minimal connectivity check: diff --git a/client/Makefile b/client/Makefile index 83887d3..51e30c1 100644 --- a/client/Makefile +++ b/client/Makefile @@ -33,7 +33,14 @@ VALKEY_LDFLAGS = -I$(VALKEY_INCLUDE_PATH) PRISKV_TARGETS_ALL = priskv-client priskv-benchmark priskv-example priskv-test_runtime PRISKV_TARGETS = priskv-client priskv-benchmark PRISKV_TARGETS_SRCS = $(patsubst priskv-%, %.c, $(PRISKV_TARGETS_ALL)) -SRCS = $(filter-out $(PRISKV_TARGETS_SRCS) $(VALKEY_BENCHMARK_SRC), $(wildcard *.c transport/*.c)) + +ifneq (,$(filter $(WITH_RDMA),yes YES y Y 1)) +CFLAGS += -DWITH_RDMA +else +RDMA_EXCLUDE = rdma.c transport/rdma.c +endif + +SRCS = $(filter-out $(PRISKV_TARGETS_SRCS) $(VALKEY_BENCHMARK_SRC) $(RDMA_EXCLUDE), $(wildcard *.c transport/*.c)) OBJS := $(SRCS:%.c=%.o) DEPS := $(OBJS:%.o=%.d) diff --git a/client/benchmark.c b/client/benchmark.c index cb3338e..7f3ec77 100644 --- a/client/benchmark.c +++ b/client/benchmark.c @@ -996,8 +996,9 @@ static void priskv_drv_get(void *ctx, const char *key, void *value, uint32_t val zctx->value = value; zctx->value_len = value_len; priskv_ctx->job->last_stage = "ACQUIRE"; - priskv_acquire_async(priskv_ctx->client, key, PRISKV_KEY_MAX_TIMEOUT, false /* pin_on_acquire */, - (uint64_t)zctx, zc_acquire_cb); + priskv_acquire_async(priskv_ctx->client, key, PRISKV_KEY_MAX_TIMEOUT, + false /* pin_on_acquire */, 0 /* pin_ttl_ms */, (uint64_t)zctx, + zc_acquire_cb); return; } /* Remove duplicate ZeroCopy GET branch */ diff --git a/client/client.c b/client/client.c index 955655a..0a2eb3b 100644 --- a/client/client.c +++ b/client/client.c @@ -275,7 +275,8 @@ static void get_handler_base(client_context *ctx, char *args, bool acquire) priskv_memory_region region = {0}; printf("ACQUIRE key=%s\n", key); /* Do not pin on acquire by default from CLI */ - status = priskv_acquire(ctx->client, key, false, 0 /* pin_ttl_ms */, ®ion); + status = priskv_acquire(ctx->client, key, PRISKV_KEY_MAX_TIMEOUT, false, 0 /* pin_ttl_ms */, + ®ion); if (status != PRISKV_STATUS_OK) { printf("Failed to GET, status(%d): %s\n", status, priskv_status_str(status)); return; @@ -477,7 +478,8 @@ static void acquire_only_handler(client_context *ctx, char *args) /* Align output field name with CLI flag semantics */ printf("ACQUIRE key=%s [PIN=%d, TTL(ms)=%" PRIu64 "]\n", key, pin_on_acquire, pin_ttl_ms); - status = priskv_acquire(ctx->client, key, pin_on_acquire, pin_ttl_ms, ®ion); + status = priskv_acquire(ctx->client, key, PRISKV_KEY_MAX_TIMEOUT, pin_on_acquire, pin_ttl_ms, + ®ion); printf("ACQUIRE status(%d): %s, addr %p, length %u, token 0x%lx\n", status, priskv_status_str(status), (void *)region.addr, region.length, region.token); if (status == PRISKV_STATUS_OK) { diff --git a/client/priskv.h b/client/priskv.h index 5e5592f..25a9cbf 100644 --- a/client/priskv.h +++ b/client/priskv.h @@ -162,9 +162,10 @@ int priskv_alloc_async(priskv_client *client, const char *key, uint32_t alloc_le int priskv_seal_async(priskv_client *client, const uint64_t *token, bool pin_on_seal, uint64_t pin_ttl_ms, uint64_t request_id, priskv_generic_cb cb); -/* Acquire memory region for zero copy read (pin_ttl_ms is PIN TTL in ms; 0 uses server default) */ -int priskv_acquire_async(priskv_client *client, const char *key, bool pin_on_acquire, - uint64_t pin_ttl_ms, uint64_t request_id, priskv_generic_cb cb); +/* Acquire: @timeout is key/op timeout for transport; @pin_ttl_ms is PIN TTL (0 = server default). */ +int priskv_acquire_async(priskv_client *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, uint64_t request_id, + priskv_generic_cb cb); /* Release memory region (by token pointer, reuse key field) */ int priskv_release_async(priskv_client *client, const uint64_t *token, bool unpin_on_release, @@ -224,9 +225,8 @@ int priskv_alloc(priskv_client *client, const char *key, uint32_t alloc_length, int priskv_seal(priskv_client *client, const uint64_t *token, bool pin_on_seal, uint64_t pin_ttl_ms); -int priskv_acquire(priskv_client *client, const char *key, bool pin_on_acquire, uint64_t pin_ttl_ms, - priskv_memory_region *region); - +int priskv_acquire(priskv_client *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, priskv_memory_region *region); int priskv_release(priskv_client *client, const uint64_t *token, bool unpin_on_release); int priskv_drop(priskv_client *client, const uint64_t *token); diff --git a/client/sync.c b/client/sync.c index bf1fd06..e87227d 100644 --- a/client/sync.c +++ b/client/sync.c @@ -190,12 +190,11 @@ int priskv_seal(priskv_client *client, const uint64_t *token, bool pin_on_seal, return req_sync.status; } -int priskv_acquire(priskv_client *client, const char *key, bool pin_on_acquire, uint64_t pin_ttl_ms, - priskv_memory_region *region) +int priskv_acquire(priskv_client *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, priskv_memory_region *region) { priskv_transport_zero_copy_req_sync req_sync = {.status = 0xffff, .done = false}; - /* pin_ttl_ms is per-request PIN TTL in ms; 0 uses server default */ - priskv_acquire_async(client, key, pin_on_acquire, pin_ttl_ms, (uint64_t)&req_sync, + priskv_acquire_async(client, key, timeout, pin_on_acquire, pin_ttl_ms, (uint64_t)&req_sync, priskv_zero_copy_req_sync_cb); priskv_sync_wait(client, &req_sync.done); diff --git a/client/transport/transport.c b/client/transport/transport.c index daad80b..58c6aab 100644 --- a/client/transport/transport.c +++ b/client/transport/transport.c @@ -37,7 +37,9 @@ priskv_transport_driver *g_client_driver = NULL; extern priskv_transport_driver priskv_transport_driver_ucx; +#ifdef WITH_RDMA extern priskv_transport_driver priskv_transport_driver_rdma; +#endif static int priskv_build_check(void) { @@ -73,10 +75,12 @@ static void __attribute__((constructor)) priskv_client_transport_init(void) driver = &priskv_transport_driver_ucx; priskv_log_notice("Using UCX transport backend\n"); break; +#ifdef WITH_RDMA case PRISKV_TRANSPORT_BACKEND_RDMA: driver = &priskv_transport_driver_rdma; priskv_log_notice("Using RDMA transport backend\n"); break; +#endif default: priskv_log_error("Unknown transport backend: %d\n", backend); break; @@ -429,17 +433,16 @@ int priskv_seal_async(priskv_client *client, const uint64_t *token, bool pin_on_ { uint32_t flags = pin_on_seal ? PRISKV_REQ_FLAG_PIN_ON_SEAL : 0; priskv_send_command(client, request_id, (const char *)token, 0 /* alloc_length */, NULL, 0, - 0 /* key_expiry_timeout */, pin_ttl_ms /* pin_ttl_ms */, - PRISKV_COMMAND_SEAL, flags, cb); + 0 /* key_expiry_timeout */, pin_ttl_ms, PRISKV_COMMAND_SEAL, flags, cb); return 0; } -int priskv_acquire_async(priskv_client *client, const char *key, bool pin_on_acquire, - uint64_t pin_ttl_ms, uint64_t request_id, priskv_generic_cb cb) +int priskv_acquire_async(priskv_client *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, uint64_t request_id, + priskv_generic_cb cb) { uint32_t flags = pin_on_acquire ? PRISKV_REQ_FLAG_PIN_ON_ACQUIRE : 0; - priskv_send_command(client, request_id, key, 0 /* alloc_length */, NULL, 0, - 0 /* key_expiry_timeout */, pin_ttl_ms /* pin_ttl_ms */, + priskv_send_command(client, request_id, key, 0 /* alloc_length */, NULL, 0, timeout, pin_ttl_ms, PRISKV_COMMAND_ACQUIRE, flags, cb); return 0; } diff --git a/client/transport/ucx.c b/client/transport/ucx.c index da328b7..12fa275 100644 --- a/client/transport/ucx.c +++ b/client/transport/ucx.c @@ -303,9 +303,11 @@ static void priskv_ucx_recv_resp_cb(ucs_status_t status, ucp_tag_t sender_tag, s return; } - if (length != sizeof(priskv_response)) { - priskv_log_warn("UCX: recv %d, expected %ld\n", length, sizeof(priskv_response)); - return; + if (ucs_unlikely(length != sizeof(priskv_response))) { + /* UCX 1.12 may provide unreliable `info->length` for tag recv callbacks. + * The response payload is a fixed-size struct, so keep going when UCS_OK. */ + priskv_log_warn("UCX: recv %zu, expected %zu (continuing)\n", length, + sizeof(priskv_response)); } uint64_t request_id = be64toh(resp->request_id); @@ -413,7 +415,6 @@ static int priskv_ucx_handshake(priskv_transport_conn *conn, ucp_address_t **add uint32_t *address_len) { int ret; - ucs_status_t status; uint8_t *peer_worker_address = NULL; size_t hs_size = sizeof(priskv_cm_ucx_handshake) + conn->worker->address_len; @@ -446,21 +447,19 @@ static int priskv_ucx_handshake(priskv_transport_conn *conn, ucp_address_t **add hs->cap.max_inflight_command = htobe16(conn->param.max_inflight_command); hs->address_len = htobe32(conn->worker->address_len); memcpy(hs->address, conn->worker->address, conn->worker->address_len); - status = ucs_socket_send(conn->connfd, hs, hs_size); + ret = priskv_safe_send(conn->connfd, hs, hs_size, NULL, NULL); free(hs); - if (status != UCS_OK) { - priskv_log_error("UCX: failed to send capability to server, status: %s\n", - ucs_status_string(status)); + if (ret) { + priskv_log_error("UCX: failed to send capability to server\n"); ret = -1; goto error; } /* receive response from server */ priskv_cm_ucx_handshake peer_hs; - status = ucs_socket_recv(conn->connfd, &peer_hs, sizeof(peer_hs)); - if (status != UCS_OK) { - priskv_log_error("UCX: failed to receive handshake msg from server, status: %s\n", - ucs_status_string(status)); + ret = priskv_safe_recv(conn->connfd, &peer_hs, sizeof(peer_hs), NULL, NULL); + if (ret) { + priskv_log_error("UCX: failed to receive handshake msg from server\n"); ret = -1; goto error; } @@ -493,10 +492,10 @@ static int priskv_ucx_handshake(priskv_transport_conn *conn, ucp_address_t **add ret = -1; goto error; } - status = ucs_socket_recv(conn->connfd, peer_worker_address, peer_worker_address_len); - if (status != UCS_OK) { - priskv_log_error("UCX: failed to receive peer_worker_address from server, status: %s\n", - ucs_status_string(status)); + ret = priskv_safe_recv(conn->connfd, peer_worker_address, peer_worker_address_len, NULL, + NULL); + if (ret) { + priskv_log_error("UCX: failed to receive peer_worker_address from server\n"); ret = -1; goto error; } diff --git a/cluster/client/Makefile b/cluster/client/Makefile index 99d8dc9..a105796 100644 --- a/cluster/client/Makefile +++ b/cluster/client/Makefile @@ -25,7 +25,7 @@ ifneq (,$(filter $(PRISKV_USE_CUDA),yes YES y Y 1)) CFLAGS += $(CUDA_LDFLAGS) -DPRISKV_USE_CUDA endif -PRISKV_CLUSTER_TARGETS = priskv-cluster-benchmark priskv-cluster-example +PRISKV_CLUSTER_TARGETS = priskv-cluster-benchmark priskv-cluster-example priskv-cluster-test_status PRISKV_CLUSTER_TARGETS_SRCS = $(patsubst priskv-cluster-%, %.c, $(PRISKV_CLUSTER_TARGETS)) ALL_SRCS = $(wildcard *.c) LIB_SRCS = $(filter-out $(PRISKV_CLUSTER_TARGETS_SRCS), $(wildcard *.c)) diff --git a/cluster/client/client.c b/cluster/client/client.c index 53891ea..1e03cbc 100644 --- a/cluster/client/client.c +++ b/cluster/client/client.c @@ -440,9 +440,8 @@ static priskvClusterRequest *priskvClusterRequestNew(priskvClusterNode *node, pr void *cbarg, RequestType type, const char *key, uint64_t timeout, priskvClusterClient *client) { - /* For ACQUIRE/SEAL, interpret 'timeout' from callers as pin_ttl_ms to keep API stable */ return priskvClusterRequestNewBase(node, sgl, nsgl, cb, NULL, cbarg, type, key, timeout, - false /*pin_key*/, 0 /*pin_ttl_ms*/, false /* unpin_key */, + false /* pin_key */, 0 /* pin_ttl_ms */, false /* unpin_key */, 0, 0, client); } @@ -452,9 +451,8 @@ priskvClusterZeroCopyRequestNew(priskvClusterNode *node, priskvClusterZeroCopyCa uint32_t alloc_length, uint64_t timeout, bool pin_key, uint64_t pin_ttl_ms, bool unpin_key, priskvClusterClient *client) { - /* Zero-copy requests never set unpin_key at construction */ return priskvClusterRequestNewBase(node, NULL, 0, NULL, cb, cbarg, type, key, timeout, pin_key, - pin_ttl_ms, unpin_key, alloc_length, 0 /*token*/, client); + pin_ttl_ms, unpin_key, alloc_length, 0 /* token */, client); } static void priskvClusterRequestFree(priskvClusterRequest *req) @@ -576,7 +574,6 @@ priskvClusterRequest *priskvClusterGetZeroCopyRequest(priskvClusterClient *clien return NULL; } - /* Use explicit pin_ttl_ms for ACQUIRE/SEAL; ALLOC ignores it */ priskvClusterRequest *req = priskvClusterZeroCopyRequestNew( node, cb, cbarg, type, key, alloc_length, timeout, pin_key, pin_ttl_ms, unpin_key, client); return req; @@ -610,18 +607,16 @@ int priskvClusterSubmitRequest(priskvClusterRequest *req) (uint64_t)req, priskvClusterZeroCopyRequestCallback); break; case SEAL: - priskv_seal_async(req->node->client, &req->token, req->pin_key /* pin_on_seal */, - req->pin_ttl_ms /* pin_ttl_ms */, (uint64_t)req, - priskvClusterRequestCallback); + priskv_seal_async(req->node->client, &req->token, req->pin_key, req->pin_ttl_ms, + (uint64_t)req, priskvClusterRequestCallback); break; case ACQUIRE: - priskv_acquire_async(req->node->client, req->key, req->pin_ttl_ms, - req->pin_key /* pin_on_acquire */, (uint64_t)req, + priskv_acquire_async(req->node->client, req->key, req->timeout, req->pin_key, + req->pin_ttl_ms, (uint64_t)req, priskvClusterZeroCopyRequestCallback); break; case RELEASE: - priskv_release_async(req->node->client, &req->token, - req->unpin_key /* unpin_on_release */, (uint64_t)req, + priskv_release_async(req->node->client, &req->token, req->unpin_key, (uint64_t)req, priskvClusterRequestCallback); break; case DROP: @@ -938,11 +933,11 @@ int priskvClusterAsyncSeal(priskvClusterClient *client, const char *key, const u bool pin_on_seal, uint64_t pin_ttl_ms, priskvClusterCallback cb, void *cbarg) { - priskvClusterRequest *req = priskvClusterGetRequest( - client, key, NULL, 0, cb, cbarg, pin_ttl_ms /* use timeout as pin_ttl_ms */, SEAL); + priskvClusterRequest *req = priskvClusterGetRequest(client, key, NULL, 0, cb, cbarg, 0, SEAL); if (req) { req->token = token ? *token : 0; req->pin_key = pin_on_seal; + req->pin_ttl_ms = pin_ttl_ms; } if (req == NULL) return -1; @@ -950,12 +945,13 @@ int priskvClusterAsyncSeal(priskvClusterClient *client, const char *key, const u return priskvClusterSubmitRequest(req); } -int priskvClusterAsyncAcquire(priskvClusterClient *client, const char *key, bool pin_on_acquire, - uint64_t pin_ttl_ms, priskvClusterZeroCopyCallback cb, void *cbarg) +int priskvClusterAsyncAcquire(priskvClusterClient *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, + priskvClusterZeroCopyCallback cb, void *cbarg) { priskvClusterRequest *req = priskvClusterGetZeroCopyRequest( - client, key, cb, cbarg, 0 /* alloc_length */, 0 /* timeout */, pin_on_acquire /* pin_key */, - pin_ttl_ms, false /* unpin_key */, ACQUIRE); + client, key, cb, cbarg, 0 /* alloc_length */, timeout, pin_on_acquire, pin_ttl_ms, + false /* unpin_key */, ACQUIRE); if (req == NULL) return -1; @@ -1103,7 +1099,8 @@ priskvClusterStatus priskvClusterAcquire(priskvClusterClient *client, const char } priskv_memory_region region = {0}; - priskv_status status = priskv_acquire(node->client, key, pin_on_acquire, pin_ttl_ms, ®ion); + priskv_status status = + priskv_acquire(node->client, key, timeout, pin_on_acquire, pin_ttl_ms, ®ion); if (status == PRISKV_STATUS_OK) { if (addr) { *addr = region.addr; @@ -1117,7 +1114,7 @@ priskvClusterStatus priskvClusterAcquire(priskvClusterClient *client, const char } int priskvClusterAcquireRegion(priskvClusterClient *client, const char *key, uint64_t timeout, - bool pin_on_acquire, uint64_t pin_ttl_ms, + bool pin_on_acquire, uint64_t pin_ttl_ms, priskv_memory_region *region) { priskvClusterNode *node = priskvClusterGetNode(client, key); @@ -1125,8 +1122,7 @@ int priskvClusterAcquireRegion(priskvClusterClient *client, const char *key, uin return PRISKV_CLUSTER_STATUS_NO_SUCH_KEY; } - /* timeout is kept for API compatibility; pin_ttl_ms controls PIN TTL (0 uses server default) */ - return priskv_acquire(node->client, key, pin_on_acquire, pin_ttl_ms, region); + return priskv_acquire(node->client, key, timeout, pin_on_acquire, pin_ttl_ms, region); } priskvClusterStatus priskvClusterRelease(priskvClusterClient *client, const char *key, diff --git a/cluster/client/client.h b/cluster/client/client.h index 3797de7..12e7397 100644 --- a/cluster/client/client.h +++ b/cluster/client/client.h @@ -89,8 +89,9 @@ int priskvClusterAsyncAlloc(priskvClusterClient *client, const char *key, uint64 int priskvClusterAsyncSeal(priskvClusterClient *client, const char *key, const uint64_t *token, bool pin_on_seal, uint64_t pin_ttl_ms, priskvClusterCallback cb, void *cbarg); -int priskvClusterAsyncAcquire(priskvClusterClient *client, const char *key, bool pin_on_acquire, - uint64_t pin_ttl_ms, priskvClusterZeroCopyCallback cb, void *cbarg); +int priskvClusterAsyncAcquire(priskvClusterClient *client, const char *key, uint64_t timeout, + bool pin_on_acquire, uint64_t pin_ttl_ms, + priskvClusterZeroCopyCallback cb, void *cbarg); int priskvClusterAsyncRelease(priskvClusterClient *client, const char *key, const uint64_t *token, bool unpin_key, priskvClusterCallback cb, void *cbarg); int priskvClusterAsyncDrop(priskvClusterClient *client, const char *key, const uint64_t *token, @@ -105,7 +106,7 @@ priskvClusterStatus priskvClusterAlloc(priskvClusterClient *client, const char * priskvClusterStatus priskvClusterSeal(priskvClusterClient *client, const char *key, const uint64_t *token, bool pin_on_seal, uint64_t pin_ttl_ms); priskvClusterStatus priskvClusterAcquire(priskvClusterClient *client, const char *key, - uint64_t timeout, bool pin_on_acquire, uint64_t pin_ttl_ms, + uint64_t timeout, bool pin_on_acquire, uint64_t pin_ttl_ms, uint64_t *addr_offset, uint32_t *valuelen); priskvClusterStatus priskvClusterRelease(priskvClusterClient *client, const char *key, const uint64_t *token, bool unpin_on_release); @@ -121,5 +122,5 @@ priskvClusterStatus priskvClusterStatusFromPriskvStatus(priskv_status status); int priskvClusterAllocRegion(priskvClusterClient *client, const char *key, uint32_t alloc_length, uint64_t timeout, priskv_memory_region *region); int priskvClusterAcquireRegion(priskvClusterClient *client, const char *key, uint64_t timeout, - bool pin_on_acquire, uint64_t pin_ttl_ms, + bool pin_on_acquire, uint64_t pin_ttl_ms, priskv_memory_region *region); diff --git a/cluster/client/test_status.c b/cluster/client/test_status.c new file mode 100644 index 0000000..9fb6d56 --- /dev/null +++ b/cluster/client/test_status.c @@ -0,0 +1,71 @@ +// Simple consistency test: verify numeric and string mappings between single-node status and cluster status +// Run: make -C cluster && ./cluster/client/priskv-cluster-test-status + +#include +#include + +#include "priskv.h" +#include "client.h" + +static int check_one(priskv_status s) +{ + priskvClusterStatus cs = priskvClusterStatusFromPriskvStatus(s); + const char *ss = priskv_status_str(s); + const char *cs_str = priskv_cluster_status_str(cs); + + int ok = 1; + + if (cs != (priskvClusterStatus)s) { + fprintf(stderr, "[NUMERIC] mismatch: priskv=%d cluster=%d\n", (int)s, (int)cs); + ok = 0; + } + + if (strcmp(ss, cs_str) != 0) { + fprintf(stderr, "[STRING] mismatch: priskv=\"%s\" cluster=\"%s\" (code=%d)\n", ss, cs_str, + (int)s); + ok = 0; + } + + return ok; +} + +int main(void) +{ + priskv_status cases[] = { + PRISKV_STATUS_OK, + PRISKV_STATUS_INVALID_COMMAND, + PRISKV_STATUS_KEY_EMPTY, + PRISKV_STATUS_KEY_TOO_BIG, + PRISKV_STATUS_VALUE_EMPTY, + PRISKV_STATUS_VALUE_TOO_BIG, + PRISKV_STATUS_NO_SUCH_COMMAND, + PRISKV_STATUS_NO_SUCH_KEY, + PRISKV_STATUS_NO_SUCH_TOKEN, + PRISKV_STATUS_INVALID_SGL, + PRISKV_STATUS_INVALID_REGEX, + PRISKV_STATUS_KEY_UPDATING, + PRISKV_STATUS_CONNECT_ERROR, + PRISKV_STATUS_SERVER_ERROR, + PRISKV_STATUS_PERMISSION_DENIED, + PRISKV_STATUS_NO_MEM, + PRISKV_STATUS_DISCONNECTED, + PRISKV_STATUS_TRANSPORT_ERROR, + PRISKV_STATUS_BUSY, + PRISKV_STATUS_PROTOCOL_ERROR, + }; + + int pass = 1; + for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); ++i) { + if (!check_one(cases[i])) { + pass = 0; + } + } + + if (!pass) { + fprintf(stderr, "Status consistency test: FAILED\n"); + return 1; + } + + printf("Status consistency test: PASSED\n"); + return 0; +} diff --git a/include/priskv-utils.h b/include/priskv-utils.h index 14e4bea..7f75a0c 100644 --- a/include/priskv-utils.h +++ b/include/priskv-utils.h @@ -148,6 +148,55 @@ static inline void priskv_inet_ntop(struct sockaddr *addr, char *dst) } } +static inline int priskv_sock_io(int sock, ssize_t (*sock_call)(int, void *, size_t, int), + int poll_events, void *data, size_t size, + void (*progress)(void *arg), void *arg, const char *name) +{ + size_t total = 0; + struct pollfd pfd; + int ret; + + while (total < size) { + pfd.fd = sock; + pfd.events = poll_events; + pfd.revents = 0; + + ret = poll(&pfd, 1, 1); /* poll for 1ms */ + if (ret > 0) { + ret = sock_call(sock, (char *)data + total, size - total, 0); + if ((ret == 0) && (poll_events & POLLIN)) { + return -1; + } + if (ret < 0) { + return -1; + } + total += ret; + } else if ((ret < 0) && (errno != EINTR)) { + return -1; + } + + /* progress user context */ + if (progress != NULL) { + progress(arg); + } + } + return 0; +} + +static inline int priskv_safe_send(int sock, void *data, size_t size, void (*progress)(void *arg), + void *arg) +{ + typedef ssize_t (*sock_call)(int, void *, size_t, int); + + return priskv_sock_io(sock, (sock_call)send, POLLOUT, data, size, progress, arg, "send"); +} + +static inline int priskv_safe_recv(int sock, void *data, size_t size, void (*progress)(void *arg), + void *arg) +{ + return priskv_sock_io(sock, recv, POLLIN, data, size, progress, arg, "recv"); +} + static inline unsigned long priskv_rdtsc(void) { unsigned long low, high; diff --git a/lib/config.c b/lib/config.c index af1e95e..8f5148f 100644 --- a/lib/config.c +++ b/lib/config.c @@ -15,6 +15,7 @@ #include "priskv-config.h" #include "priskv-log.h" #include +#include static const char *priskv_transport_backend_names[] = {[PRISKV_TRANSPORT_BACKEND_RDMA] = "RDMA", [PRISKV_TRANSPORT_BACKEND_UCX] = "UCX", @@ -152,7 +153,8 @@ void priskv_config_init(void) ucs_status_t priskv_config_parser_fill_opts(void *opts, ucs_config_global_list_entry_t *entry, const char *env_prefix, int ignore_errors) { - return ucs_config_parser_fill_opts(opts, entry, env_prefix, ignore_errors); + /* UCX v1.12+: ucs_config_parser_fill_opts(opts, fields, env_prefix, table_prefix, ignore_errors) */ + return ucs_config_parser_fill_opts(opts, entry->table, env_prefix, entry->prefix, ignore_errors); } void priskv_config_parser_release_opts(void *opts, ucs_config_field_t *fields) @@ -163,6 +165,6 @@ void priskv_config_parser_release_opts(void *opts, ucs_config_field_t *fields) ucs_status_t priskv_config_parser_set_value(void *opts, ucs_config_field_t *fields, const char *prefix, const char *name, const char *value) { - - return ucs_config_parser_set_value(opts, fields, prefix, name, value); + (void)prefix; + return ucs_config_parser_set_value(opts, fields, name, value); } diff --git a/lib/ucx.c b/lib/ucx.c index 59781de..080e6b1 100644 --- a/lib/ucx.c +++ b/lib/ucx.c @@ -205,8 +205,8 @@ ucs_status_t priskv_ucx_munmap(priskv_ucx_memh *memh) } if (memh->rkey_buffer) { - ucp_memh_buffer_release_params_t params = {.field_mask = 0}; - ucp_memh_buffer_release(memh->rkey_buffer, ¶ms); + /* Packed RKEY buffer release (UCX 1.12+). */ + ucp_rkey_buffer_release(memh->rkey_buffer); memh->rkey_buffer = NULL; } ucp_mem_unmap(memh->context->handle, memh->handle); @@ -255,7 +255,10 @@ priskv_ucx_worker *priskv_ucx_worker_create(priskv_ucx_context *context, uint64_ worker->efd = -1; } - status = ucp_worker_get_address(worker->handle, &worker->address, &worker->address_len); + /* UCX uses size_t* for address length; wire format uses uint32_t. */ + size_t addr_len = 0; + status = ucp_worker_get_address(worker->handle, &worker->address, &addr_len); + worker->address_len = (uint32_t)addr_len; PRISKV_UCX_RETURN_IF_ERROR( status, "priskv_ucx_worker_init: failed to get address", { free(worker); }, NULL); @@ -497,7 +500,8 @@ static ucs_status_ptr_t priskv_ucx_post_tag_recv(priskv_ucx_worker *worker, pris ucp_request_param_t param = {.op_attr_mask = UCP_OP_ATTR_FIELD_CALLBACK | UCP_OP_ATTR_FIELD_DATATYPE | - UCP_OP_ATTR_FIELD_USER_DATA, + UCP_OP_ATTR_FIELD_USER_DATA | + UCP_OP_ATTR_FIELD_RECV_INFO, .datatype = ucp_dt_make_contig(1), .cb.recv = priskv_ucx_request_tag_recv_cb_intl, .user_data = request}; diff --git a/pypriskv/priskv/priskv_client.py b/pypriskv/priskv/priskv_client.py index bd0157a..56bb93a 100644 --- a/pypriskv/priskv/priskv_client.py +++ b/pypriskv/priskv/priskv_client.py @@ -83,11 +83,12 @@ def seal(self, def acquire(self, key: str, - pin_ttl_ms: int = 0, + timeout: int = client.PRISKV_KEY_MAX_TIMEOUT, pin_on_acquire: bool = False, + pin_ttl_ms: int = 0, ) -> Tuple[int, client.MemoryRegion]: - """Acquire zero-copy read region. When pin_on_acquire is True, pin_ttl_ms specifies the PIN TTL in milliseconds; 0 means server default.""" - status, region = client.acquire(self.conn, key, pin_on_acquire, pin_ttl_ms) + """Acquire: timeout is transport/key timeout; pin_ttl_ms is PIN TTL when pin_on_acquire is True.""" + status, region = client.acquire(self.conn, key, timeout, pin_on_acquire, pin_ttl_ms) return status, region def release(self, diff --git a/pypriskv/pybind.cpp b/pypriskv/pybind.cpp index f0a6576..9b7b58b 100644 --- a/pypriskv/pybind.cpp +++ b/pypriskv/pybind.cpp @@ -139,12 +139,13 @@ std::tuple priskv_alloc_wrapper(uintptr_t client, std::string key } std::tuple priskv_acquire_wrapper(uintptr_t client, std::string key, - uint64_t pin_ttl_ms, bool pin_on_acquire) + uint64_t timeout, bool pin_on_acquire, + uint64_t pin_ttl_ms) { uint64_t addr_offset = 0; uint32_t value_length = 0; - int ret = priskvClusterAcquire((priskvClusterClient *)client, key.c_str(), 0 /* timeout */, + int ret = priskvClusterAcquire((priskvClusterClient *)client, key.c_str(), timeout, pin_on_acquire, pin_ttl_ms, &addr_offset, &value_length); return {ret, addr_offset, value_length}; } @@ -327,18 +328,17 @@ PYBIND11_MODULE(_priskv_client, m) m.def( "acquire", - [](uintptr_t client, std::string key, bool pin_on_acquire, uint64_t pin_ttl_ms) { + [](uintptr_t client, std::string key, uint64_t timeout, bool pin_on_acquire, + uint64_t pin_ttl_ms) { priskv_memory_region region {0}; - /* Pass timeout=0 and explicit pin_ttl_ms */ - int ret = - priskvClusterAcquireRegion((priskvClusterClient *)client, key.c_str(), - 0 /* timeout */, pin_on_acquire, pin_ttl_ms, ®ion); + int ret = priskvClusterAcquireRegion((priskvClusterClient *)client, key.c_str(), timeout, + pin_on_acquire, pin_ttl_ms, ®ion); return py::make_tuple(ret, region); }, - py::arg("client"), py::arg("key"), py::arg("pin_on_acquire") = false, - py::arg("pin_ttl_ms") = 0, - "A function to acquire memory region for read; when PIN is requested, pin_ttl_ms sets TTL " - "in ms (0 uses server default)."); + py::arg("client"), py::arg("key"), py::arg("timeout") = PRISKV_KEY_MAX_TIMEOUT, + py::arg("pin_on_acquire") = false, py::arg("pin_ttl_ms") = 0, + "A function to acquire memory region for read; timeout is transport/key timeout; " + "pin_ttl_ms is PIN TTL when pin_on_acquire is true (0 uses server default)."); m.def( "release", diff --git a/pypriskv/testing.py b/pypriskv/testing.py index c05a1ab..77c660c 100644 --- a/pypriskv/testing.py +++ b/pypriskv/testing.py @@ -540,17 +540,18 @@ def test_unpin_no_such_key(self): def test_pin_ttl_expire_on_acquire(self): TEST_KEY = "py_pin_ttl_acquire" SIZE = 256 + TIMEOUT = 3000 # TTL is in milliseconds; sleep slightly longer than 1.5s to avoid clock granularity issues PIN_TTL_MS = 1500 SLEEP_SEC = 1.6 - status, region = self.client.alloc(TEST_KEY, SIZE, 3000) + status, region = self.client.alloc(TEST_KEY, SIZE, TIMEOUT) assert status == 0 status = self.client.seal(TEST_KEY, region) assert status == 0 - # ACQUIRE with PIN(with TTL) - status, acq = self.client.acquire(TEST_KEY, PIN_TTL_MS, True) + # ACQUIRE with PIN and explicit PIN TTL (transport timeout is TIMEOUT) + status, acq = self.client.acquire(TEST_KEY, TIMEOUT, pin_on_acquire=True, pin_ttl_ms=PIN_TTL_MS) assert status == 0, f"acquire(pin ttl) failed: {status}" # wait TTL expired @@ -560,7 +561,7 @@ def test_pin_ttl_expire_on_acquire(self): assert status == priskv.PRISKV_STATUS.PRISKV_STATUS_UNPIN_NOT_CLOSED, \ f"expected UNPIN_NOT_CLOSED after TTL expiry, got {status}" - status2, acq2 = self.client.acquire(TEST_KEY, 0, False) + status2, acq2 = self.client.acquire(TEST_KEY, TIMEOUT, pin_on_acquire=False, pin_ttl_ms=0) assert status2 == priskv.PRISKV_STATUS.PRISKV_STATUS_OK status2 = self.client.release(TEST_KEY, acq2, unpin_on_release=True) assert status2 == priskv.PRISKV_STATUS.PRISKV_STATUS_UNPIN_NOT_CLOSED, \ diff --git a/run_e2e_test.py b/run_e2e_test.py index 88024de..f1f7f79 100755 --- a/run_e2e_test.py +++ b/run_e2e_test.py @@ -28,7 +28,11 @@ def find_rdma_dev(): for dev in os.listdir(ibclass): netdev = ibclass + dev + "/ports/1/gid_attrs/ndevs/0" with open(netdev) as fp: - addrs = netifaces.ifaddresses(fp.readline().strip("\n")) + iface = fp.readline().strip("\n") + try: + addrs = netifaces.ifaddresses(iface) + except ValueError: + continue if netifaces.AF_INET in addrs: ipv4_addr = addrs[netifaces.AF_INET][0]["addr"] print( diff --git a/run_unit_test.py b/run_unit_test.py index 3168117..d009a6b 100755 --- a/run_unit_test.py +++ b/run_unit_test.py @@ -35,7 +35,7 @@ def priskv_unit_test(parallel: bool = True): "./server/test/test-buddy-mt", "./server/test/test-kv", "./server/test/test-kv-mt", "./server/test/test-memory --no-tmpfs", "./server/test/test-slab", "./server/test/test-transport 8 200", - "./server/test/test-kv-pin-ttl" + "./server/test/test-kv-pin-ttl", ] print("---- PrisKV UNIT TEST ----") diff --git a/server/Makefile b/server/Makefile index 4e5c471..bd5322e 100644 --- a/server/Makefile +++ b/server/Makefile @@ -38,12 +38,18 @@ else CFLAGS += -O0 endif +ifneq (,$(filter $(WITH_RDMA),yes YES y Y 1)) +CFLAGS += -DWITH_RDMA +else +RDMA_EXCLUDE = ./rdma.c ./transport/rdma.c +endif + PRISKV_BINPATH = usr/bin PRISKV_MANPATH = usr/share/man PRISKV_SERVER_TARGETS = priskv-server priskv-memfile PRISKV_SERVER_TARGETS_SRCS = $(patsubst priskv-%, ./%.c, $(PRISKV_SERVER_TARGETS)) -PRISKV_SERVER_SRCS = $(filter-out $(PRISKV_SERVER_TARGETS_SRCS), $(shell find . -path ./test -prune -o -name "*.c" -print)) +PRISKV_SERVER_SRCS = $(filter-out $(PRISKV_SERVER_TARGETS_SRCS) $(RDMA_EXCLUDE), $(shell find . -path ./test -prune -o -name "*.c" -print)) PRISKV_SERVER_OBJS := $(PRISKV_SERVER_SRCS:%.c=%.o) PRISKV_SERVER_DEPS := $(PRISKV_SERVER_OBJS:%.o=%.d) diff --git a/server/kv.c b/server/kv.c index b79376f..71e00c3 100644 --- a/server/kv.c +++ b/server/kv.c @@ -107,6 +107,64 @@ typedef struct priskv_kv { } pin_stats; } priskv_kv; +/* + * TODO(wangyi): Implement PinTTL cleanup mechanism + * + * Context: + * - Pin operations (PIN_ON_ACQUIRE, PIN_ON_SEAL) increase pin_count on the latest visible + * version to protect keys from eviction. If a consumer crashes or a request fails, some + * pin operations may never be closed by UNPIN (e.g., RELEASE with UNPIN_ON_RELEASE), leaving + * keys indefinitely pinned. + * + * Goals: + * - Introduce a best-effort TTL-based cleanup for orphaned pins so that keys are eventually + * unpinned when their associated requests are gone. + * - Maintain multi-version correctness: UNPIN targets the latest version even if the pin was + * created before a SEAL publish that replaced the visible version. + * + * Proposed design: + * - PinOperator: a lightweight record created on each effective pin, containing: + * - key (bytes + length) + * - creation timestamp (monotonic clock) + * - ttl_ms (configurable per pin or global default) + * - optional origin (ACQUIRE or SEAL) and a debug request_id for observability + * - PinManager: per-bucket or global manager storing PinOperator entries in an expiry-ordered + * min-heap or timing-wheel to enable O(logN) insert and efficient batch expiry checks. + * - On pin: + * - After pin_count++ on the targeted keynode, create and register a PinOperator. + * - On unpin: + * - Remove the corresponding PinOperator (match by key); then decrement pin_count on latest + * version using priskv_key_unpin_latest semantics. Multiple pins on the same key will + * have multiple PinOperator entries. + * - On seal (version migration): + * - pin_count is already inherited to the new version; PinOperator records keep referencing + * the key (not the keynode pointer), so no migration is required. + * - Scheduling: + * - Reuse expire routine infrastructure (timerfd) to periodically check PinManager and + * perform cleanup for expired entries. Each expired PinOperator triggers + * priskv_key_unpin_latest(kv, old_keynode_of_record) on the latest version by key. + * - Consider sharding the PinManager by hash-bucket index to minimize global contention. + * - Concurrency & locking: + * - PinOperator insert/remove should use per-bucket spinlocks consistent with hash-head + * protection, avoiding deadlocks by keeping lock order (manager lock -> keynode lock). + * - Configuration: + * - Provide a global default TTL (e.g., kv->pin_ttl_ms) and allow per-request override via + * request flags or auxiliary fields in the protocol header (future extension). + * - Observability & safeguards: + * - Export counters: pin_ttl_active, pin_ttl_expired, pin_ttl_cleanup_ops, pin_ttl_orphaned. + * - Cap the maximum number of active PinOperator entries to prevent memory blow-up; when + * exceeding the cap, log warnings and fallback to immediate unpin or refuse new pins. + * - Recovery: + * - PinTTL metadata is best-effort and may be non-persistent. After restart, keys might + * remain pinned by pin_count; scheduled cleanup resumes with new PinOperator records for + * future pins. Persistent logging can be considered if stronger guarantees are needed. + * + * Next steps: + * - Add PinManager data structures and lifecycle APIs. + * - Integrate with pin/unpin paths and expire routine scheduling. + * - Extend info endpoints to expose PinTTL metrics. + */ + static void priskv_lru_access(priskv_key *keynode, bool is_in_list) { priskv_kv *kv = keynode->kv; @@ -937,6 +995,23 @@ int priskv_drop_node(void *_kv, void *_keynode) return PRISKV_RESP_STATUS_OK; } +/* Increment pin_count on the given keynode. */ +int priskv_key_pin(void *_kv, void *_keynode) +{ + priskv_kv *kv = (priskv_kv *)_kv; + priskv_key *keynode = (priskv_key *)_keynode; + if (!keynode) { + return PRISKV_RESP_STATUS_SERVER_ERROR; + } + pthread_spin_lock(&keynode->lock); + keynode->pin_count++; + pthread_spin_unlock(&keynode->lock); + if (kv) { + kv->pin_stats.pin_ops++; + } + return PRISKV_RESP_STATUS_OK; +} + /* Decrement pin_count on the latest version of the key corresponding to keynode. */ int priskv_key_unpin_latest(void *_kv, void *_keynode) { @@ -954,10 +1029,10 @@ int priskv_key_unpin_latest(void *_kv, void *_keynode) priskv_resp_status resp = priskv_pin_count_delta_latest(kv, node->key, node->keylen, -1, 0); /* Stats: count every UNPIN attempt; record not-closed cases. */ - __sync_fetch_and_add(&kv->pin_stats.unpin_ops, 1); - if (resp == PRISKV_RESP_STATUS_UNPIN_NOT_CLOSED) { - __sync_fetch_and_add(&kv->pin_stats.unpin_not_closed, 1); - } + __sync_fetch_and_add(&kv->pin_stats.unpin_ops, 1); + if (resp == PRISKV_RESP_STATUS_UNPIN_NOT_CLOSED) { + __sync_fetch_and_add(&kv->pin_stats.unpin_not_closed, 1); + } if (resp != PRISKV_RESP_STATUS_OK) { priskv_log_warn("KV: UNPIN_LATEST status=%d for node=%p (len=%u)\n", resp, (void *)node, node->keylen); diff --git a/server/kv.h b/server/kv.h index 061e654..41ed15d 100644 --- a/server/kv.h +++ b/server/kv.h @@ -165,6 +165,8 @@ int priskv_publish_node(void *_kv, void *_keynode); int priskv_publish_node_with_pin(void *_kv, void *_keynode, bool pin_on_publish, uint64_t ttl_ms); int priskv_drop_node(void *_kv, void *_keynode); +/* Pin/Unpin controls for lifecycle protection */ +int priskv_key_pin(void *_kv, void *_keynode); int priskv_key_unpin_latest(void *_kv, void *_keynode); /* Pin on the latest version; ttl_ms==0 uses the default TTL */ int priskv_key_pin_latest(void *_kv, void *_keynode, uint64_t ttl_ms); @@ -174,7 +176,6 @@ uint64_t priskv_get_pin_ops(void *_kv); uint64_t priskv_get_pin_failed_ops(void *_kv); uint64_t priskv_get_unpin_ops(void *_kv); uint64_t priskv_get_unpin_not_closed(void *_kv); - #if defined(__cplusplus) } #endif diff --git a/server/test/Makefile b/server/test/Makefile index ce8a819..27cad45 100644 --- a/server/test/Makefile +++ b/server/test/Makefile @@ -33,7 +33,7 @@ ifneq (,$(filter $(PRISKV_USE_CUDA),yes YES y Y 1)) CFLAGS += $(CUDA_LDFLAGS) -DPRISKV_USE_CUDA endif -.PHONY: $(TEST_BUDDY) ${TEST_BUDDY_MT} $(TEST_SLAB) $(TEST_SLAB_MT) $(TEST_KV) $(TEST_KV_MT) $(TEST_TRANSPORT) $(TEST_MEMORY) $(TEST_ACL) $(TEST_KV_EXPIRE_ROUTINE) $(TEST_BE_REDIS) +.PHONY: $(TEST_BUDDY) ${TEST_BUDDY_MT} $(TEST_SLAB) $(TEST_SLAB_MT) $(TEST_KV) $(TEST_KV_MT) $(TEST_TRANSPORT) $(TEST_MEMORY) $(TEST_ACL) $(TEST_KV_EXPIRE_ROUTINE) $(TEST_BE_REDIS) $(TEST_KV_PIN_TTL) OBJS = ../memory.o ../kv.o ../slab.o ../crc.o ../acl.o all: $(TEST_BUDDY) ${TEST_BUDDY_MT} $(TEST_SLAB) $(TEST_SLAB_MT) $(TEST_KV) $(TEST_KV_MT) $(TEST_TRANSPORT) $(TEST_MEMORY) $(TEST_ACL) $(TEST_KV_EXPIRE_ROUTINE) $(TEST_BE_REDIS) $(TEST_KV_PIN_TTL) @@ -87,7 +87,6 @@ else $(TEST_KV_EXPIRE_ROUTINE): @echo "UCX not available, skip $(TEST_KV_EXPIRE_ROUTINE)" endif - ifeq ($(shell pkg-config --exists ucx; echo $$?),0) $(TEST_KV_PIN_TTL): $(OBJS) $(CC) test_kv_pin_ttl.c ../../lib/workqueue.c ../../lib/threads.c ../../lib/event.c ../../lib/ucx.c ../../lib/config.c ../memory.c ../kv.c ../slab.c ../buddy.c ../crc.c ../../lib/log.c ../backend/backend.c ../transport/transport.c ../transport/rdma.c ../transport/ucx.c ../tiering.c ../acl.c $(CFLAGS) $(UCX_CFLAGS) -o $(TEST_KV_PIN_TTL) -lmount -lpthread -lrdmacm -libverbs $(UCX_LIBS) @@ -96,7 +95,6 @@ $(TEST_KV_PIN_TTL): @echo "UCX not available, skip $(TEST_KV_PIN_TTL)" endif - $(TEST_BE_REDIS): $(CC) test_be_redis.c ../../lib/log.c ../../lib/event.c ../../lib/workqueue.c ../../lib/threads.c ../backend/backend.c ../backend/be_redis.c $(CFLAGS) -o $(TEST_BE_REDIS) -levent -lhiredis diff --git a/server/transport/transport.c b/server/transport/transport.c index 8edb376..ee8f176 100644 --- a/server/transport/transport.c +++ b/server/transport/transport.c @@ -38,7 +38,9 @@ priskv_transport_server g_transport_server = { bool priskv_test_token_add_fail_once = false; extern priskv_transport_driver priskv_transport_driver_ucx; +#ifdef WITH_RDMA extern priskv_transport_driver priskv_transport_driver_rdma; +#endif uint32_t g_slow_query_threshold_latency_us = SLOW_QUERY_THRESHOLD_LATENCY_US; @@ -60,10 +62,12 @@ static void __attribute__((constructor)) priskv_server_transport_init(void) driver = &priskv_transport_driver_ucx; priskv_log_notice("Using UCX transport backend\n"); break; +#ifdef WITH_RDMA case PRISKV_TRANSPORT_BACKEND_RDMA: driver = &priskv_transport_driver_rdma; priskv_log_notice("Using RDMA transport backend\n"); break; +#endif default: priskv_log_error("Unknown transport backend: %d\n", backend); break; @@ -310,18 +314,35 @@ int priskv_transport_handle_recv(priskv_transport_conn *conn, priskv_request *re return -EPROTO; } - keylen = len - keyoff; + /* UCX callback's `length` may be unreliable under some UCX versions. + * Rely on the protocol header's key_length instead of deriving it from `len`. */ + keylen = be16toh(req->key_length); + if (!keylen) { - priskv_log_warn("Transport: <%s - %s> empty key. recv %d, less than %d, nsgl 0x%x\n", + priskv_log_warn("Transport: <%s - %s> empty key. len(%u) keyoff(%u) nsgl 0x%x\n", conn->local_addr, conn->peer_addr, len, keyoff, nsgl); driver->send_response(conn, req->request_id, PRISKV_RESP_STATUS_KEY_EMPTY, 0, 0, 0); return -EPROTO; } + /* Optional sanity check: if UCX `len` is sane and smaller than what we expect, treat as protocol error. */ + if (len != 0 && len < (uint32_t)(keyoff + keylen)) { + priskv_log_warn("Transport: <%s - %s> invalid key. recv len(%u) < keyoff(%u)+keylen(%u)\n", + conn->local_addr, conn->peer_addr, len, keyoff, keylen); + driver->send_response(conn, req->request_id, PRISKV_RESP_STATUS_INVALID_COMMAND, 0, 0, 0); + return -EPROTO; + } + if (keylen > conn->conn_cap.max_key_length) { - priskv_log_warn("Transport: <%s - %s> invalid key. key(%d) exceeds max_key_length(%d)\n", - conn->local_addr, conn->peer_addr, keylen, conn->conn_cap.max_key_length); + uint16_t raw_key_length = be16toh(req->key_length); + uint16_t raw_nsgl = be16toh(req->nsgl); + uint32_t raw_alloc_length = be32toh(req->alloc_length); + priskv_log_warn( + "Transport: <%s - %s> invalid key. len(%u) keyoff(%u) keylen(%u) " + "key_length(%u) nsgl(%u) alloc_length(%u) exceeds max_key_length(%u)\n", + conn->local_addr, conn->peer_addr, len, keyoff, keylen, raw_key_length, raw_nsgl, + raw_alloc_length, conn->conn_cap.max_key_length); driver->send_response(conn, req->request_id, PRISKV_RESP_STATUS_KEY_TOO_BIG, 0, 0, 0); return -EPROTO; } @@ -593,14 +614,10 @@ int priskv_transport_handle_recv(priskv_transport_conn *conn, priskv_request *re PRISKV_RESP_STATUS_PERMISSION_DENIED, 0, 0, 0); break; } - /* Atomically publish and optionally pin; treat req.timeout as the TTL (ms) for this - * pin. */ - /* Use pin_ttl_ms; 0 means default TTL. Only read TTL when pin_on_seal is set. */ + /* Atomically publish and optionally pin; pin_ttl_ms 0 means default TTL. */ status = priskv_publish_node_with_pin( conn->kv, keynode, (flags & PRISKV_REQ_FLAG_PIN_ON_SEAL) != 0, (flags & PRISKV_REQ_FLAG_PIN_ON_SEAL) ? pin_ttl_ms : 0); - /* TTL registration is handled in the KV layer (publish critical section); no need to - * duplicate here. */ priskv_transport_token_del(conn, token); ret = driver->send_response(conn, req->request_id, status, 0, 0, 0); } @@ -614,26 +631,18 @@ int priskv_transport_handle_recv(priskv_transport_conn *conn, priskv_request *re } { uint64_t addr_offset = 0; - /* Enforce ACQUIRE+PIN all-or-nothing semantics when requested: - * - If PIN_ON_ACQUIRE is set, attempt to pin the latest visible version first. - * - If pin fails (e.g., key deleted concurrently), fail the whole ACQUIRE and release - * the reference acquired by priskv_get_key. - */ + /* Enforce ACQUIRE+PIN all-or-nothing semantics when requested. */ if (flags & PRISKV_REQ_FLAG_PIN_ON_ACQUIRE) { - /* Use pin_ttl_ms; 0 means default TTL. */ priskv_resp_status presp = priskv_key_pin_latest(conn->kv, keynode, pin_ttl_ms); if (presp != PRISKV_RESP_STATUS_OK) { - /* Atomicity: no token created, drop our reference and return pin's status */ priskv_get_key_end(keynode); ret = driver->send_response(conn, req->request_id, presp, 0, 0, 0); break; } } - /* Create token only after (optional) pin succeeds to keep ACQUIRE+PIN atomic */ uint64_t token = priskv_transport_token_add(conn, keynode, PRISKV_TOKEN_TYPE_ACQUIRE); if (!token) { - /* Roll back pin if we pinned, then release the key reference */ if (flags & PRISKV_REQ_FLAG_PIN_ON_ACQUIRE) { (void)priskv_key_unpin_latest(conn->kv, keynode); } @@ -642,7 +651,6 @@ int priskv_transport_handle_recv(priskv_transport_conn *conn, priskv_request *re 0, 0, 0); break; } - /* TTL registration is handled in the KV layer (pin path); no need to duplicate here. */ status = priskv_value_addr_offset(conn->kv, val, &addr_offset); ret = driver->send_response(conn, req->request_id, status, valuelen, addr_offset, token); @@ -675,19 +683,10 @@ int priskv_transport_handle_recv(priskv_transport_conn *conn, priskv_request *re PRISKV_RESP_STATUS_PERMISSION_DENIED, 0, 0, 0); break; } - /* RELEASE semantics when UNPIN is requested: - * - If PRISKV_REQ_FLAG_UNPIN_ON_RELEASE is set and unpin fails (e.g., NO_SUCH_KEY / - * UNPIN_NOT_CLOSED), we currently release the ACQUIRE reference and delete the - * token as a best-effort cleanup, and return the unpin status to the client. - * - Otherwise, we release the reference and delete the token, returning OK. - * Note: This favors simplicity over retry semantics; adjust if caller requires a - * strict all-or-nothing behavior. - */ priskv_resp_status resp = PRISKV_RESP_STATUS_OK; if (flags & PRISKV_REQ_FLAG_UNPIN_ON_RELEASE) { resp = priskv_key_unpin_latest(conn->kv, keynode); } - /* Either UNPIN succeeded or not requested: finish RELEASE */ priskv_get_key_end(keynode); priskv_transport_token_del(conn, token); ret = driver->send_response(conn, req->request_id, resp, 0, 0, 0); diff --git a/server/transport/ucx.c b/server/transport/ucx.c index 8758e26..b516916 100644 --- a/server/transport/ucx.c +++ b/server/transport/ucx.c @@ -237,10 +237,9 @@ static inline void priskv_ucx_reject(priskv_transport_conn *client, priskv_cm_st .status = htobe16(status), .value = htobe64(value), }; - ucs_status_t ucs_status = ucs_socket_send(client->connfd, &rej_msg_be, sizeof(rej_msg_be)); - if (ucs_status != UCS_OK) { - priskv_log_error("UCX: send reject message failed, status: %s\n", - ucs_status_string(ucs_status)); + int ret = priskv_safe_send(client->connfd, &rej_msg_be, sizeof(rej_msg_be), NULL, NULL); + if (ret < 0) { + priskv_log_error("UCX: send reject message failed: %m\n"); } } @@ -275,11 +274,9 @@ static inline int priskv_ucx_accept(priskv_transport_conn *client) client->peer_addr, address_len, print_len, worker_address_hex); } - ucs_status_t status = ucs_socket_send(client->connfd, hs, hs_size); - if (status != UCS_OK) { - ret = -1; - priskv_log_error("UCX: send accept message failed, status: %s\n", - ucs_status_string(status)); + ret = priskv_safe_send(client->connfd, hs, hs_size, NULL, NULL); + if (ret < 0) { + priskv_log_error("UCX: send accept message failed: %m\n"); goto out_free_msg; } @@ -294,7 +291,6 @@ static inline int priskv_ucx_accept(priskv_transport_conn *client) static inline int priskv_ucx_handle_handshake(void *arg) { int ret; - ucs_status_t sock_status; priskv_cm_ucx_handshake peer_hs; priskv_cm_status status; @@ -304,10 +300,9 @@ static inline int priskv_ucx_handle_handshake(void *arg) int connfd = client->connfd; /* #step0, recv handshake msg */ - sock_status = ucs_socket_recv(connfd, &peer_hs, sizeof(peer_hs)); - if (sock_status != UCS_OK) { - priskv_log_error("UCX: recv handshake msg failed, status: %s\n", - ucs_status_string(sock_status)); + ret = priskv_safe_recv(connfd, &peer_hs, sizeof(peer_hs), NULL, NULL); + if (ret < 0) { + priskv_log_error("UCX: recv handshake msg failed: %m\n"); ucs_close_fd(&connfd); return -1; } @@ -325,10 +320,9 @@ static inline int priskv_ucx_handle_handshake(void *arg) ucs_close_fd(&connfd); return -1; } - sock_status = ucs_socket_recv(connfd, peer_worker_address, peer_worker_address_len); - if (sock_status != UCS_OK) { - priskv_log_error("UCX: recv peer address failed, status: %s\n", - ucs_status_string(sock_status)); + ret = priskv_safe_recv(connfd, peer_worker_address, peer_worker_address_len, NULL, NULL); + if (ret < 0) { + priskv_log_error("UCX: recv peer address failed: %m\n"); ucs_close_fd(&connfd); return -1; } @@ -478,19 +472,15 @@ static inline void priskv_ucx_handle_cm(int fd, void *opaque, uint32_t ev) connfd = accept(listener->listenfd, (struct sockaddr *)&client_addr, &client_addr_len); if (connfd < 0) { if (errno == EAGAIN || errno == EWOULDBLOCK) { - // no connection available, return return; } else if (errno == EINTR) { - // interrupted by signal, try again goto again; } else { - // other errors, log and return priskv_log_error("UCX: accept on listenfd %d failed: %m\n", listener->listenfd); return; } } - // got a connection priskv_inet_ntop(&client_addr, peer_addr); priskv_log_info("UCX: accept on listenfd %d, connfd %d, client addr %s\n", listener->listenfd, connfd, peer_addr); From 491b38d2052f720f89257556852072f62ecc4d6a Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Wed, 18 Mar 2026 23:13:22 +0000 Subject: [PATCH 3/7] fix: remove .pyc from tracking, add __pycache__ to gitignore Signed-off-by: bob-021206 --- .gitignore | 3 +++ 1 file changed, 3 insertions(+) diff --git a/.gitignore b/.gitignore index 0e42b57..1160218 100644 --- a/.gitignore +++ b/.gitignore @@ -31,6 +31,9 @@ pypriskv/priskv.egg-info/ pypriskv/priskv/*.so include/priskv-version.h pypriskv/priskv/__pycache__/ +pypriskv/**/__pycache__/ +**/__pycache__/ +*.pyc man/priskv-server.1 .vscode/ tmp/ From 2a67d3b74e9cbc22930f1df0da652ad40ba89779 Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Thu, 19 Mar 2026 19:16:37 +0000 Subject: [PATCH 4/7] docs: add libucx-dev to prerequisites in setup guide Signed-off-by: bob-021206 --- SETUP_UCX_TCP.md | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/SETUP_UCX_TCP.md b/SETUP_UCX_TCP.md index 0404074..f5be66b 100644 --- a/SETUP_UCX_TCP.md +++ b/SETUP_UCX_TCP.md @@ -23,7 +23,8 @@ apt-get install -y \ dpkg-dev debhelper \ pkg-config \ python3-pybind11 python3-dev python3-pip \ - libonig-dev libhiredis-dev liburing-dev + libonig-dev libhiredis-dev liburing-dev \ + libucx-dev libucx0 ``` ### Python build tooling @@ -85,7 +86,7 @@ pip install -U pip setuptools wheel cd /PrisKV make all cd pypriskv -pip install -v -e . +pip install --no-build-isolation -v -e . ``` ### Option B: Install from wheel @@ -127,6 +128,7 @@ export PRISKV_LOG_LEVEL=notice ## 5. Start `priskv-server` (UCX TCP) ```bash +source .venv/bin/activate cd /PrisKV export PRISKV_TRANSPORT=ucx export UCX_TLS=tcp From d528190345717c9f95beb540e884534c6182bb88 Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Wed, 22 Apr 2026 21:34:40 +0000 Subject: [PATCH 5/7] fix: align UCS config API with system UCX (no flags field, print_opts arity) Signed-off-by: bob-021206 --- include/priskv-config.h | 3 +-- lib/config.c | 8 ++++---- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/include/priskv-config.h b/include/priskv-config.h index cc0e0ad..ebb30d3 100644 --- a/include/priskv-config.h +++ b/include/priskv-config.h @@ -42,8 +42,7 @@ extern "C" .prefix = PREFIX, \ .table = TABLE, \ .size = sizeof(TYPE), \ - .list = {NULL, NULL}, \ - .flags = 0}; + .list = {NULL, NULL}}; #define PRISKV_CONFIG_GET_TABLE(TABLE) &g_##TABLE##_config_entry #define PRISKV_ENV_PREFIX "PRISKV_" diff --git a/lib/config.c b/lib/config.c index 8f5148f..835422d 100644 --- a/lib/config.c +++ b/lib/config.c @@ -129,10 +129,10 @@ static void priskv_config_init_impl(void) "Shared memory will be enabled automatically."); } - ucs_config_parser_print_opts( - stdout, "PrisKV Environment Variables", &g_config, priskv_config_table, NULL, - PRISKV_ENV_PREFIX, - UCS_CONFIG_PRINT_CONFIG | UCS_CONFIG_PRINT_HEADER | UCS_CONFIG_PRINT_DOC, NULL); + ucs_config_parser_print_opts(stdout, "PrisKV Environment Variables", &g_config, + priskv_config_table, NULL, PRISKV_ENV_PREFIX, + UCS_CONFIG_PRINT_CONFIG | UCS_CONFIG_PRINT_HEADER | + UCS_CONFIG_PRINT_DOC); // logging priskv_set_log_level(g_config.logging.log_level); From 68d1de5079877f2b7ec24bc9d20d75e8408ec390 Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Sat, 9 May 2026 18:29:52 -0400 Subject: [PATCH 6/7] Remove unused code and update docs --- SETUP_UCX_TCP.md | 2 +- cluster/client/test_status.c | 71 ---------------------------------- server/kv.c | 75 ------------------------------------ server/kv.h | 1 - 4 files changed, 1 insertion(+), 148 deletions(-) delete mode 100644 cluster/client/test_status.c diff --git a/SETUP_UCX_TCP.md b/SETUP_UCX_TCP.md index f5be66b..57b584e 100644 --- a/SETUP_UCX_TCP.md +++ b/SETUP_UCX_TCP.md @@ -120,7 +120,7 @@ export PRISKV_CLIENT_DIRECT_MODE=y Optional debugging: ```bash -export PRISKV_LOG_LEVEL=notice +export PRISKV_LOG_LEVEL=debug ``` --- diff --git a/cluster/client/test_status.c b/cluster/client/test_status.c deleted file mode 100644 index 9fb6d56..0000000 --- a/cluster/client/test_status.c +++ /dev/null @@ -1,71 +0,0 @@ -// Simple consistency test: verify numeric and string mappings between single-node status and cluster status -// Run: make -C cluster && ./cluster/client/priskv-cluster-test-status - -#include -#include - -#include "priskv.h" -#include "client.h" - -static int check_one(priskv_status s) -{ - priskvClusterStatus cs = priskvClusterStatusFromPriskvStatus(s); - const char *ss = priskv_status_str(s); - const char *cs_str = priskv_cluster_status_str(cs); - - int ok = 1; - - if (cs != (priskvClusterStatus)s) { - fprintf(stderr, "[NUMERIC] mismatch: priskv=%d cluster=%d\n", (int)s, (int)cs); - ok = 0; - } - - if (strcmp(ss, cs_str) != 0) { - fprintf(stderr, "[STRING] mismatch: priskv=\"%s\" cluster=\"%s\" (code=%d)\n", ss, cs_str, - (int)s); - ok = 0; - } - - return ok; -} - -int main(void) -{ - priskv_status cases[] = { - PRISKV_STATUS_OK, - PRISKV_STATUS_INVALID_COMMAND, - PRISKV_STATUS_KEY_EMPTY, - PRISKV_STATUS_KEY_TOO_BIG, - PRISKV_STATUS_VALUE_EMPTY, - PRISKV_STATUS_VALUE_TOO_BIG, - PRISKV_STATUS_NO_SUCH_COMMAND, - PRISKV_STATUS_NO_SUCH_KEY, - PRISKV_STATUS_NO_SUCH_TOKEN, - PRISKV_STATUS_INVALID_SGL, - PRISKV_STATUS_INVALID_REGEX, - PRISKV_STATUS_KEY_UPDATING, - PRISKV_STATUS_CONNECT_ERROR, - PRISKV_STATUS_SERVER_ERROR, - PRISKV_STATUS_PERMISSION_DENIED, - PRISKV_STATUS_NO_MEM, - PRISKV_STATUS_DISCONNECTED, - PRISKV_STATUS_TRANSPORT_ERROR, - PRISKV_STATUS_BUSY, - PRISKV_STATUS_PROTOCOL_ERROR, - }; - - int pass = 1; - for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); ++i) { - if (!check_one(cases[i])) { - pass = 0; - } - } - - if (!pass) { - fprintf(stderr, "Status consistency test: FAILED\n"); - return 1; - } - - printf("Status consistency test: PASSED\n"); - return 0; -} diff --git a/server/kv.c b/server/kv.c index 71e00c3..6a68d8e 100644 --- a/server/kv.c +++ b/server/kv.c @@ -107,64 +107,6 @@ typedef struct priskv_kv { } pin_stats; } priskv_kv; -/* - * TODO(wangyi): Implement PinTTL cleanup mechanism - * - * Context: - * - Pin operations (PIN_ON_ACQUIRE, PIN_ON_SEAL) increase pin_count on the latest visible - * version to protect keys from eviction. If a consumer crashes or a request fails, some - * pin operations may never be closed by UNPIN (e.g., RELEASE with UNPIN_ON_RELEASE), leaving - * keys indefinitely pinned. - * - * Goals: - * - Introduce a best-effort TTL-based cleanup for orphaned pins so that keys are eventually - * unpinned when their associated requests are gone. - * - Maintain multi-version correctness: UNPIN targets the latest version even if the pin was - * created before a SEAL publish that replaced the visible version. - * - * Proposed design: - * - PinOperator: a lightweight record created on each effective pin, containing: - * - key (bytes + length) - * - creation timestamp (monotonic clock) - * - ttl_ms (configurable per pin or global default) - * - optional origin (ACQUIRE or SEAL) and a debug request_id for observability - * - PinManager: per-bucket or global manager storing PinOperator entries in an expiry-ordered - * min-heap or timing-wheel to enable O(logN) insert and efficient batch expiry checks. - * - On pin: - * - After pin_count++ on the targeted keynode, create and register a PinOperator. - * - On unpin: - * - Remove the corresponding PinOperator (match by key); then decrement pin_count on latest - * version using priskv_key_unpin_latest semantics. Multiple pins on the same key will - * have multiple PinOperator entries. - * - On seal (version migration): - * - pin_count is already inherited to the new version; PinOperator records keep referencing - * the key (not the keynode pointer), so no migration is required. - * - Scheduling: - * - Reuse expire routine infrastructure (timerfd) to periodically check PinManager and - * perform cleanup for expired entries. Each expired PinOperator triggers - * priskv_key_unpin_latest(kv, old_keynode_of_record) on the latest version by key. - * - Consider sharding the PinManager by hash-bucket index to minimize global contention. - * - Concurrency & locking: - * - PinOperator insert/remove should use per-bucket spinlocks consistent with hash-head - * protection, avoiding deadlocks by keeping lock order (manager lock -> keynode lock). - * - Configuration: - * - Provide a global default TTL (e.g., kv->pin_ttl_ms) and allow per-request override via - * request flags or auxiliary fields in the protocol header (future extension). - * - Observability & safeguards: - * - Export counters: pin_ttl_active, pin_ttl_expired, pin_ttl_cleanup_ops, pin_ttl_orphaned. - * - Cap the maximum number of active PinOperator entries to prevent memory blow-up; when - * exceeding the cap, log warnings and fallback to immediate unpin or refuse new pins. - * - Recovery: - * - PinTTL metadata is best-effort and may be non-persistent. After restart, keys might - * remain pinned by pin_count; scheduled cleanup resumes with new PinOperator records for - * future pins. Persistent logging can be considered if stronger guarantees are needed. - * - * Next steps: - * - Add PinManager data structures and lifecycle APIs. - * - Integrate with pin/unpin paths and expire routine scheduling. - * - Extend info endpoints to expose PinTTL metrics. - */ - static void priskv_lru_access(priskv_key *keynode, bool is_in_list) { priskv_kv *kv = keynode->kv; @@ -995,23 +937,6 @@ int priskv_drop_node(void *_kv, void *_keynode) return PRISKV_RESP_STATUS_OK; } -/* Increment pin_count on the given keynode. */ -int priskv_key_pin(void *_kv, void *_keynode) -{ - priskv_kv *kv = (priskv_kv *)_kv; - priskv_key *keynode = (priskv_key *)_keynode; - if (!keynode) { - return PRISKV_RESP_STATUS_SERVER_ERROR; - } - pthread_spin_lock(&keynode->lock); - keynode->pin_count++; - pthread_spin_unlock(&keynode->lock); - if (kv) { - kv->pin_stats.pin_ops++; - } - return PRISKV_RESP_STATUS_OK; -} - /* Decrement pin_count on the latest version of the key corresponding to keynode. */ int priskv_key_unpin_latest(void *_kv, void *_keynode) { diff --git a/server/kv.h b/server/kv.h index 41ed15d..9e9ab02 100644 --- a/server/kv.h +++ b/server/kv.h @@ -166,7 +166,6 @@ int priskv_publish_node_with_pin(void *_kv, void *_keynode, bool pin_on_publish, int priskv_drop_node(void *_kv, void *_keynode); /* Pin/Unpin controls for lifecycle protection */ -int priskv_key_pin(void *_kv, void *_keynode); int priskv_key_unpin_latest(void *_kv, void *_keynode); /* Pin on the latest version; ttl_ms==0 uses the default TTL */ int priskv_key_pin_latest(void *_kv, void *_keynode, uint64_t ttl_ms); From ea9637c5025788da54bd8e498ec0f923268f7f8d Mon Sep 17 00:00:00 2001 From: bob-021206 Date: Sat, 9 May 2026 18:39:20 -0400 Subject: [PATCH 7/7] delete outdated logs in SETUP_UCX_TCP.md --- SETUP_UCX_TCP.md | 32 +------------------------------- 1 file changed, 1 insertion(+), 31 deletions(-) diff --git a/SETUP_UCX_TCP.md b/SETUP_UCX_TCP.md index 57b584e..aeac072 100644 --- a/SETUP_UCX_TCP.md +++ b/SETUP_UCX_TCP.md @@ -6,8 +6,6 @@ The goal is: 2. Connect using `PriskvClient` to a running `priskv-server` 3. Optionally validate end-to-end `set/get` -> Note: PrisKV may require UCX-version-specific compatibility patches. If your system uses UCX 1.12, verify the items in **Section 7** exist in your source tree. - --- ## 0. Prerequisites (Debian/Ubuntu) @@ -183,35 +181,7 @@ PY --- -## 7. UCX 1.12 compatibility patches to verify - -If your system UCX is 1.12, PrisKV may need compatibility adjustments to avoid build errors or protocol/ABI mismatches. -Verify the following items exist in your checkout. - -### 7.1 `PrisKV/lib/config.c`: UCX config parser API compatibility - -- Add `` (avoid implicit `strcmp`) -- Fix `ucs_config_parser_print_opts` call signature (argument count) -- Fix wrapper signatures/arguments for `ucs_config_parser_fill_opts` and `ucs_config_parser_set_value` - -### 7.2 `PrisKV/include/priskv-config.h`: config table struct compatibility - -- In `PRISKV_CONFIG_DECLARE_TABLE`, remove `.flags = 0` if the field does not exist in UCX 1.12. - -### 7.3 `PrisKV/lib/ucx.c`: UCX 1.12 API compatibility - -- Replace packed RKEY release logic with UCX 1.12 behavior (`ucp_rkey_buffer_release()`) -- Adjust `ucp_worker_get_address()` length argument to `size_t*` and convert back to PrisKV types as needed - -### 7.4 `PrisKV/lib/ucx.c` (client/server UCX wrapper): tag-recv completion info - -- In `priskv_ucx_post_tag_recv()`, ensure `ucp_request_param_t.op_attr_mask` includes `UCP_OP_ATTR_FIELD_RECV_INFO` - -If you still see messages like `UCX: recv <...>, expected 48` or endpoint timeouts, continue deeper protocol-level debugging (request/response completion info and response struct decoding). - ---- - -## 8. Readiness checklist for developers +## 7. Readiness checklist for developers - Successful `import priskv` indicates the Python extension is built/installed correctly. - Server logs show `UCX ... ready` indicates UCX transport is initialized.