Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,10 @@ target_include_directories(kvspace
$<INSTALL_INTERFACE:include>
)

# 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)

Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,两个后端与前端三者共用同一份契约。
Expand Down
11 changes: 9 additions & 2 deletions include/kvspace/kvspace.h
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,10 @@ typedef struct {

/* ── 生命周期 ─────────────────────────────────────────────────── */
void *kvspaceConnect(const char *dsn);

/* 关闭:关闭前落盘未决写(失败写 stderr,绝不静默)。
* 进程正常退出(exit / main 返回)时,前端对所有未 Close 句柄兜底关闭即落盘——
* 漏调 Close 不丢数据,真丢也绝不静默。 */
void kvspaceClose(void *h);

/* ── 单点读写 / 目录 ──────────────────────────────────────────── */
Expand Down Expand Up @@ -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);
Expand Down
53 changes: 46 additions & 7 deletions src/frontend.c
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -13,6 +13,7 @@
#include "kvspace/kvspace.h"

#include <dlfcn.h>
#include <pthread.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
Expand Down Expand Up @@ -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。 */
Expand Down Expand Up @@ -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) ─────────────────────────────────── */
Expand Down
85 changes: 84 additions & 1 deletion tutorial-durable/test.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,23 @@
#!/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/
WS = ROOT.parent # array2d/
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 提取预期输出行。"""
Expand Down Expand Up @@ -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(
Expand All @@ -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)
Expand Down
Loading