diff --git a/CMakeLists.txt b/CMakeLists.txt index 65ec11d..33472e2 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -26,8 +26,10 @@ target_include_directories(kvspace $ ) -# dlopen/dlsym 的宿主库:Linux 为 libdl,macOS 在 libSystem 内(CMAKE_DL_LIBS 两者皆正确) -target_link_libraries(kvspace PRIVATE ${CMAKE_DL_LIBS}) +# 活跃句柄表(退出兜底落盘)用 pthread 互斥;dlopen/dlsym 的宿主库 Linux 为 libdl, +# macOS 在 libSystem 内(CMAKE_DL_LIBS 两者皆正确)。 +find_package(Threads REQUIRED) +target_link_libraries(kvspace PRIVATE ${CMAKE_DL_LIBS} Threads::Threads) target_compile_definitions(kvspace PRIVATE _GNU_SOURCE) diff --git a/README.md b/README.md index 70af8ef..ba851e2 100644 --- a/README.md +++ b/README.md @@ -7,7 +7,7 @@ KVSpace 的 dispatch 前端。消费者只链接 `libkvspace`,运行期按 DSN | `shm://...` | kvspace-c(file-backed mmap + ART) | `libkvspace-c.so.1` | | `redis://` `fs://` `s3://` | kvspace-durable(Rust) | `libkvspace_durable.so.1` | -两个后端导出同一套 25 个 `kvspace*` C ABI(见 `include/kvspace/kvspace.h`)。 +两个后端导出同一套 `kvspace*` C ABI(见 `include/kvspace/kvspace.h`)。 前端用 `dlopen(RTLD_NOW | RTLD_LOCAL)` 装载后端,handle 包一层 vtable; codec(`kvspaceTlvEncode*`/`kvspaceDecodeHead`/`kvspaceNew*`)无 handle,由前端静态实现, head 格式 byte-identical,两个后端与前端三者共用同一份契约。 diff --git a/include/kvspace/kvspace.h b/include/kvspace/kvspace.h index 30a8dee..6bc838e 100644 --- a/include/kvspace/kvspace.h +++ b/include/kvspace/kvspace.h @@ -60,6 +60,10 @@ typedef struct { /* ── 生命周期 ─────────────────────────────────────────────────── */ void *kvspaceConnect(const char *dsn); + +/* 关闭:关闭前落盘未决写(失败写 stderr,绝不静默)。 + * 进程正常退出(exit / main 返回)时,前端对所有未 Close 句柄兜底关闭即落盘—— + * 漏调 Close 不丢数据,真丢也绝不静默。 */ void kvspaceClose(void *h); /* ── 单点读写 / 目录 ──────────────────────────────────────────── */ @@ -104,13 +108,16 @@ int kvspaceSetPart(void *h, const char *key, uint32_t offset, const uint8_t *buf int kvspaceGetHead(void *h, const char *key, kvspaceHead_t *out); /* 就地写:key 必须已存在、kind 不变、body_len 必须等于原 body_len——返回原 box 的 body - * 偏移指针供调用方直接写。违反前置条件 → 非 0 + err(绝不静默重分配、绝不回落)。写即持久。 */ + * 偏移指针供调用方直接写。违反前置条件 → 非 0 + err(绝不静默重分配、绝不回落)。 + * 落盘时机:body 由调用方在返回后填,故本笔在下一次 kvspace* 调用(读也算)或 kvspaceClose + * 时落盘;前者失败以返回码 + err 报出,后者失败写 stderr。 */ int kvspaceWriteInPlace(void *h, const char *key, int resolve, uint32_t body_len, uint8_t **body, char *err, uint32_t err_cap); /* 新位置写:按 (ref, storetype, ro, vid, langtype, body_len) 分配新 box、写好 head,返回 body * 偏移指针供直接写。ARRAYND 的 dims 由 codec 从 langtype 串内的 [dims] 解析落入物理字段。 - * 用于新建 key 或 storetype/尺寸变化。写即持久。 */ + * 用于新建 key 或 storetype/尺寸变化。落盘时机同 WriteInPlace:本笔在下一次 kvspace* 调用 + * 或 kvspaceClose 时落盘。 */ int kvspaceWriteNewPlace(void *h, const char *key, uint8_t ref, uint8_t storetype, uint8_t ro, uint32_t vid, const char *langtype, uint32_t body_len, uint8_t **body, char *err, uint32_t err_cap); diff --git a/src/frontend.c b/src/frontend.c index 81bc4e2..7f51ce3 100644 --- a/src/frontend.c +++ b/src/frontend.c @@ -1,8 +1,8 @@ /* frontend.c — KVSpace dispatch 前端。 * - * 导出与两个后端完全相同的 27 个 kvspace* 符号。运行期按 DSN scheme 用 dlopen + * 导出 kvspace* C ABI(与两个后端同名同签名)。运行期按 DSN scheme 用 dlopen * (RTLD_NOW | RTLD_LOCAL)装载后端,把 handle 包一层 {vtable, dl, backend}。 - * codec(TlvEncode、DecodeHead、New 等)无 handle,由前端静态实现,byte-identical。 + * codec(TlvEncode、DecodeHead、New 等)与 kvspaceConst 无 handle,由前端静态实现,byte-identical。 * * 后端装载名(后缀随平台): * shm://... → libkvspace-c.so.1 / macOS: libkvspace-c.dylib @@ -13,6 +13,7 @@ #include "kvspace/kvspace.h" #include +#include #include #include #include @@ -56,14 +57,48 @@ typedef struct { char *err, uint32_t err_cap); } kvspace_vt; -typedef struct { +typedef struct kvspace_handle { kvspace_vt *vt; void *dl; void *backend; + struct kvspace_handle *next; } kvspace_handle; static kvspace_handle *H(void *h) { return (kvspace_handle *)h; } +/* 活跃句柄表:durable 的写侧惰性(body 由调用方在返回后填,落盘延到下一次 op),故漏调 + * Close 就退出会丢最后一笔写。进程正常退出时前端兜底落盘再关闭,兑现 kvspace.h 的退出保证。 */ +static kvspace_handle *live; +static int live_hooked; +static pthread_mutex_t live_lock = PTHREAD_MUTEX_INITIALIZER; + +static void handle_free(kvspace_handle *x) { + x->vt->free(x->backend); + dlclose(x->dl); + free(x->vt); + free(x); +} + +static void live_unlink(kvspace_handle *x) { + pthread_mutex_lock(&live_lock); + for (kvspace_handle **p = &live; *p; p = &(*p)->next) + if (*p == x) { *p = x->next; break; } + pthread_mutex_unlock(&live_lock); +} + +/* 退出兜底:未 Close 句柄逐个关闭——durable 的 Close 落盘未决写,落盘失败由后端自己写 + * stderr。数据可以丢(后端故障),但绝不静默。 */ +static void exit_close(void) { + for (;;) { + pthread_mutex_lock(&live_lock); + kvspace_handle *x = live; + if (x) live = x->next; + pthread_mutex_unlock(&live_lock); + if (!x) return; + handle_free(x); + } +} + /* ── 后端选择 ───────────────────────────────────────────────────────── */ /* 后端库名:Linux 为 .so.1,macOS 为 .dylib。 */ @@ -158,16 +193,20 @@ void *kvspaceConnect(const char *dsn) { kvspace_handle *h = malloc(sizeof(*h)); if (!h) { free(vt); dlclose(dl); return NULL; } h->vt = vt; h->dl = dl; h->backend = backend; + + pthread_mutex_lock(&live_lock); + h->next = live; + live = h; + if (!live_hooked) { live_hooked = 1; atexit(exit_close); } + pthread_mutex_unlock(&live_lock); return h; } void kvspaceClose(void *h) { if (!h) return; kvspace_handle *x = H(h); - if (x->vt->free) x->vt->free(x->backend); - if (x->dl) dlclose(x->dl); - free(x->vt); - free(x); + live_unlink(x); + handle_free(x); } /* ── 单点读写 / 目录(trampoline) ─────────────────────────────────── */ diff --git a/tutorial-durable/test.py b/tutorial-durable/test.py index da85f1e..eb7f75e 100644 --- a/tutorial-durable/test.py +++ b/tutorial-durable/test.py @@ -1,13 +1,15 @@ #!/usr/bin/env python3 # -*- coding: utf-8 -*- """运行 tutorial/*.sh,对比脚本头部注释中的 expected 输出; -并交叉校验 kvspace-c 与 kvspace-durable 的 head 编解码(rw/vid)字节一致。""" +交叉校验 kvspace-c 与 kvspace-durable 的 head 编解码(rw/vid)字节一致; +并回归写侧落盘时机:下一次 op 落盘与退出兜底(issue #23)。""" import ctypes import os import struct import subprocess import sys +import tempfile from pathlib import Path ROOT = Path(__file__).resolve().parent.parent # kvspace/ @@ -15,6 +17,7 @@ DURABLE_SO = WS / "kvspace-durable" / "target" / "release" / "libkvspace_durable.so" KVSPACE_C_DIR = WS / "kvspace-c" KVSPACE_C_SO = KVSPACE_C_DIR / "build" / "libkvspace-c.so" +FRONTEND_SO = os.environ.get("KVSPACE_SO", "/usr/lib/kvspace/libkvspace.so.1") def extract_expected(script): """从脚本头部 # expected: ... # /end 提取预期输出行。""" @@ -161,6 +164,85 @@ def test_kvspace_c_alignment(): return ok +# ── 写侧落盘时机(issue #23)───────────────────────────────────────── +# +# 写把 TLV 攒在句柄内(body 由调用方在返回后填,返回前无法落盘):下一次 op 落盘上一笔 +# (读也是 op),另一个句柄随即可见;漏调 Close 而正常退出时由前端退出兜底落盘。子进程必须 +# 独立——它正是「不 Close 就退出」的那一个。 + +CHILD = r''' +import ctypes, sys +U8P = ctypes.POINTER(ctypes.c_uint8) +lib = ctypes.CDLL(sys.argv[1]) +lib.kvspaceConnect.restype = ctypes.c_void_p +lib.kvspaceConnect.argtypes = [ctypes.c_char_p] +lib.kvspaceDelTree.argtypes = [ctypes.c_void_p, ctypes.c_char_p, ctypes.c_char_p, ctypes.c_uint32] +lib.kvspaceWriteNewPlace.argtypes = [ + ctypes.c_void_p, ctypes.c_char_p, ctypes.c_uint8, ctypes.c_uint8, ctypes.c_uint8, + ctypes.c_uint32, ctypes.c_char_p, ctypes.c_uint32, + ctypes.POINTER(U8P), ctypes.c_char_p, ctypes.c_uint32] +lib.kvspaceGet.argtypes = [ctypes.c_void_p, ctypes.c_char_p, ctypes.c_int, + ctypes.POINTER(U8P), ctypes.POINTER(ctypes.c_uint32)] +err = ctypes.create_string_buffer(256) + +def put(h, key, val): + body = U8P() + lt = ("[%d]char/utf8" % len(val)).encode() + assert lib.kvspaceWriteNewPlace(h, key, 0, 2, 0, 0, lt, len(val), + ctypes.byref(body), err, 256) == 0, err.value + ctypes.memmove(body, val, len(val)) + +def get(h, key): + out, n = U8P(), ctypes.c_uint32() + lib.kvspaceGet(h, key, 0, ctypes.byref(out), ctypes.byref(n)) + return None if not out else ctypes.string_at(out, n.value) + +h1 = lib.kvspaceConnect(sys.argv[2].encode()) +lib.kvspaceDelTree(h1, b"/flushprobe/", err, 256) +put(h1, b"/flushprobe/a", b"va") +put(h1, b"/flushprobe/b", b"vb") +put(h1, b"/flushprobe/c", b"vc") +get(h1, b"/flushprobe/a") # 任意一次 op 先把上一笔(c)落盘 +h2 = lib.kvspaceConnect(sys.argv[2].encode()) +print("next-op-visible", get(h2, b"/flushprobe/c") is not None) +put(h1, b"/flushprobe/d", b"vd") # 留作未决写:不再发 op、不 Close,直接正常退出 +sys.exit(0) +''' + + +def test_exit_flush(): + """写侧落盘时机:下一次 op 即落盘上一笔;漏调 Close 正常退出不丢最后一笔。""" + if not Path(FRONTEND_SO).exists(): + print(f'FAIL exit-flush (缺前端 {FRONTEND_SO})') + return False + kvbin = os.path.expanduser('~/.local/bin/kvspace') + env = os.environ.copy() + keys = ['/flushprobe/a', '/flushprobe/b', '/flushprobe/c', '/flushprobe/d'] + with tempfile.TemporaryDirectory() as tmp: + cases = [ + ('redis', env.get('KVSPACE', 'redis://127.0.0.1:6379')), + ('fs', f'fs://{tmp}/fs'), + ('shm', f'shm://{tmp}/probe.shm'), + ] + ok = True + for label, dsn in cases: + child = subprocess.run([sys.executable, '-c', CHILD, FRONTEND_SO, dsn], + capture_output=True, text=True, timeout=30) + seen = 'next-op-visible True' in child.stdout + # 子进程已退出(未 Close):经 CLI 复核四笔全在,最后一笔由退出兜底落盘。 + got = subprocess.run([kvbin, 'get', *keys], capture_output=True, text=True, + timeout=10, env={**env, 'KVSPACE': dsn}) + kept = [l for l in got.stdout.splitlines() if 'char/utf8:v' in l] + if seen and len(kept) == len(keys): + print(f'PASS exit-flush {label}') + else: + print(f'FAIL exit-flush {label} (next-op-visible={seen}, 落地 {len(kept)}/{len(keys)})') + if child.stderr: + print(f' child stderr: {child.stderr.strip()}') + ok = False + return ok + + def main(): here = os.path.dirname(os.path.abspath(__file__)) scripts = sorted( @@ -173,6 +255,7 @@ def main(): results = [test_script(s) for s in scripts] results.append(test_kvspace_c_alignment()) + results.append(test_exit_flush()) passed = sum(results) print(f'\n{passed}/{len(results)} passed') sys.exit(0 if passed == len(results) else 1)