diff --git a/.github/actions/setup-kvspace-c/action.yml b/.github/actions/setup-kvspace-c/action.yml new file mode 100644 index 00000000..b7e346a3 --- /dev/null +++ b/.github/actions/setup-kvspace-c/action.yml @@ -0,0 +1,35 @@ +name: setup-kvspace-c +runs: + using: composite + steps: + - uses: actions/checkout@v4 + with: + repository: array2d/blockmalloc + path: blockmalloc + - uses: actions/checkout@v4 + with: + repository: array2d/slotsboxmalloc + path: slotsboxmalloc + - uses: actions/checkout@v4 + with: + repository: ${{ github.repository_owner }}/kvspace-c + ref: feat/notify-take + path: kvspace-c + - name: Build kvspace-c + shell: bash + run: | + cmake -S blockmalloc -B blockmalloc/build -DCMAKE_BUILD_TYPE=Release + cmake --build blockmalloc/build -j --target blockmalloc + mkdir -p slotsboxmalloc/build kvspace-c/build + gcc -shared -fPIC -O2 -o slotsboxmalloc/build/libslotsboxmalloc.so \ + slotsboxmalloc/src/slotsboxmalloc.c \ + -Islotsboxmalloc/include -Iblockmalloc/include \ + -Lblockmalloc/build -lblockmalloc \ + -Wl,-rpath,$GITHUB_WORKSPACE/blockmalloc/build + gcc -shared -fPIC -O2 -o kvspace-c/build/libkvspace-c.so \ + kvspace-c/src/xvalue.c kvspace-c/src/kvspace.c kvspace-c/src/durable_abi.c \ + -Ikvspace-c/include -Iblockmalloc/include -Islotsboxmalloc/include \ + -Lblockmalloc/build -lblockmalloc -Lslotsboxmalloc/build -lslotsboxmalloc \ + -lpthread \ + -Wl,-rpath,$GITHUB_WORKSPACE/blockmalloc/build \ + -Wl,-rpath,$GITHUB_WORKSPACE/slotsboxmalloc/build diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index bb81c9bc..1fc928e7 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -9,68 +9,47 @@ on: workflow_dispatch: jobs: - # ── 编译 + 静态检查 ────────────────────────────────────────────────────── build: - strategy: - matrix: - os: [ubuntu-latest, macos-latest] - runs-on: ${{ matrix.os }} + runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 - - uses: actions/setup-go@v5 with: - go-version: "1.24" - - run: go mod tidy - - run: go build ./... + path: kvlang + - uses: ./kvlang/.github/actions/setup-kvspace-c + - name: Runtime tests + working-directory: kvlang + run: | + cmake -S runtime -B build/runtime -DCMAKE_BUILD_TYPE=Release -DKVSPACE_LIB=kvspace-c + cmake --build build/runtime --target kvlang_runtime \ + kvlang_delegate_test kvlang_notify_test kvlang_incr_test \ + kvlang_expire_test kvlang_watchany_test -j + ./bin/kvlang_delegate_test + ./bin/kvlang_notify_test + ./bin/kvlang_incr_test + ./bin/kvlang_expire_test + ./bin/kvlang_watchany_test - # ── tutorial 集成测试(需要 Redis)──────────────────────────────────────── tutorial-test: runs-on: ubuntu-latest - services: - redis: - image: redis:7-alpine - ports: - - 6379:6379 - options: >- - --health-cmd "redis-cli ping" - --health-interval 10s - --health-timeout 5s - --health-retries 5 steps: - uses: actions/checkout@v4 - - uses: actions/setup-go@v5 with: - go-version: "1.24" - - run: go mod tidy - - run: go build -o kvlang ./cmd/kvlang/ - - run: go install github.com/array2d/kvspace-go/cmd/kvspace@latest - - run: python3 -m unittest tutorial.test.BenchmarkTest - - run: python3 tutorial/test.py + path: kvlang + - uses: ./kvlang/.github/actions/setup-kvspace-c + - uses: dtolnay/rust-toolchain@stable + - name: make shm + working-directory: kvlang + run: mkdir -p bin && make shm -j + - name: tutorial + working-directory: kvlang + env: + KVSPACE: shm:///tmp/kvlang_ci + LD_LIBRARY_PATH: ${{ github.workspace }}/kvspace-c/build:${{ github.workspace }}/blockmalloc/build:${{ github.workspace }}/slotsboxmalloc/build + run: python3 tutorial/test.py --no-build - # ── 多平台交叉编译(仅 tag 时生成 release 产物)───────────────────────── release: if: startsWith(github.ref, 'refs/tags/v') runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 - - uses: actions/setup-go@v5 - with: - go-version: "1.24" - - - name: Cross-compile - run: | - go mod tidy - mkdir -p dist - for GOOS in linux darwin; do - for GOARCH in amd64 arm64; do - echo "Building $GOOS/$GOARCH..." - GOOS=$GOOS GOARCH=$GOARCH go build -ldflags="-s -w" -o dist/kvlang-$GOOS-$GOARCH ./cmd/kvlang/ - done - done - ls -lh dist/ - - - name: Create Release - uses: softprops/action-gh-release@v2 - with: - files: dist/* - generate_release_notes: true + - run: echo "0.2.0 release artifacts are built with make shm; Go dist is retired" diff --git a/layout/src/code.rs b/layout/src/code.rs index 34aa132a..752679b8 100644 --- a/layout/src/code.rs +++ b/layout/src/code.rs @@ -34,11 +34,11 @@ pub fn compile(kv: &mut Kv, src: &str) -> Result<(), String> { for func in &file.funcs { let pkg = if func.pkg.is_empty() { file.package.clone() } else { func.pkg.clone() }; let mut lowered = lower::lower_func(func); - write_func(kv, &pkg, &mut lowered); + write_func(kv, &pkg, &mut lowered)?; any_code = true; } for decl in &file.rwir_decls { - write_rwir_decl(kv, decl); + write_rwir_decl(kv, decl)?; } let mut body = file.init_body.clone(); @@ -53,7 +53,7 @@ pub fn compile(kv: &mut Kv, src: &str) -> Result<(), String> { pkg: String::new(), }; let mut lowered = lower::lower_func(&init_fn); - write_func(kv, "", &mut lowered); + write_func(kv, "", &mut lowered)?; any_code = true; } @@ -100,10 +100,21 @@ pub fn vet(src: &str) -> Result<(), String> { } /// 写函数到 /lib/:签名(rwfunc)、源码、参数 Ptr、指令体。 -pub fn write_func(kv: &mut Kv, pkg: &str, fn_: &mut Func) { +pub fn write_func(kv: &mut Kv, pkg: &str, fn_: &mut Func) -> Result<(), String> { let mut type_map = lower::infer_types(fn_); lower::specialize(fn_, &type_map); let func_dir = keytree::lib_func(pkg, &fn_.sig.name); + let opcode = if pkg.is_empty() { + fn_.sig.name.clone() + } else { + format!("{}{}{}", pkg, keytree::MEMBER_SEP, fn_.sig.name) + }; + let existing = kv.get_one(&keytree::rwir(&opcode)); + if kvkind::kind(&existing) == kvkind::KIND_RWIR { + return Err(format!( + "{opcode}: rwir declaration and rwfunc body both define it; one opcode, one definition" + )); + } // 按函数覆盖(文件夹复制式合并):只 del_tree 本函数子树,不动 /lib 下其它函数。 // 禁止整库删除——layoutcode 必须可增量:多次 layout 各自覆盖其函数,不误删先前的函数。 @@ -134,16 +145,23 @@ pub fn write_func(kv: &mut Kv, pkg: &str, fn_: &mut Func) { let _ = kv.set(&pairs); write_body(kv, pkg, &fn_.sig.name, &fn_.body, &mut type_map, 1); + Ok(()) } /// 写用户声明的 rwir(无体)到 /lib/。 -pub fn write_rwir_decl(kv: &mut Kv, decl: &RwirDecl) { +pub fn write_rwir_decl(kv: &mut Kv, decl: &RwirDecl) -> Result<(), String> { let mut opcode = decl.sig.name.clone(); if !decl.pkg.is_empty() { opcode = format!("{}{}{opcode}", decl.pkg, keytree::MEMBER_SEP); } + let sig = kv.get_one(&format!("{}/[0,0]", keytree::lib_func("", &opcode))); + if kvkind::kind(&sig) == kvkind::KIND_RWFUNC { + return Err(format!( + "{opcode}: rwir declaration and rwfunc body both define it; one opcode, one definition" + )); + } let v = kvkind::new_rwir(decl.sig.num_reads(), decl.sig.num_writes(), &decl.sig.kindexp_list().join("\n")); - let _ = kv.set(&[(keytree::rwir(&opcode), v)]); + kv.set(&[(keytree::rwir(&opcode), v)]) } /// 将 body 写入 /lib/// 下。offset 起始 idx(顶层函数=1)。 diff --git a/layout/tests/pipeline_test.rs b/layout/tests/pipeline_test.rs index 96289a0b..27a9c1f9 100644 --- a/layout/tests/pipeline_test.rs +++ b/layout/tests/pipeline_test.rs @@ -17,6 +17,14 @@ fn fresh_kv() -> Kv { kv } +#[test] +fn refuse_rwir_and_rwfunc() { + let mut kv = fresh_kv(); + compile(&mut kv, "rwfunc dup() -> () {\n 1 -> _\n}\n").unwrap(); + let err = compile(&mut kv, "rwir dup() -> ()\n").unwrap_err(); + assert!(err.contains("one opcode, one definition"), "{err}"); +} + #[test] fn compile_simple_func() { let mut kv = fresh_kv(); diff --git a/runtime/CMakeLists.txt b/runtime/CMakeLists.txt index 86805104..27cd8e1b 100644 --- a/runtime/CMakeLists.txt +++ b/runtime/CMakeLists.txt @@ -23,7 +23,7 @@ set(CMAKE_RUNTIME_OUTPUT_DIRECTORY "${BIN_DIR}") add_library(kvlang_runtime SHARED src/strbuf.c src/xvalue.c src/kv.c src/keytree.c src/rwir.c src/vthread.c src/logx.c src/builtin.c src/kvcpu.c src/runtime.c src/rwirext.c - src/type_expr.c) + src/type_expr.c src/dispatch.c) target_include_directories(kvlang_runtime PUBLIC include src) target_compile_definitions(kvlang_runtime PRIVATE _GNU_SOURCE) @@ -41,3 +41,37 @@ target_link_libraries(kvlang_runtime PRIVATE m pthread) # --export-dynamic 供 kvlang(term 扩展)找符号;--disable-new-dtags 使 rpath 转 DT_RPATH(传递), # 让 runtime → kvspace-c → blockmalloc/slotsboxmalloc 的传递依赖能被 runtime 的 rpath 找到。 set_target_properties(kvlang_runtime PROPERTIES LINK_FLAGS "-Wl,--export-dynamic -Wl,--disable-new-dtags") + +add_executable(kvlang_delegate_test tests/delegate_test.c) +target_link_libraries(kvlang_delegate_test PRIVATE kvlang_runtime pthread) +target_include_directories(kvlang_delegate_test PRIVATE include src) +add_executable(kvlang_notify_test tests/notify_contract_test.c) +target_link_libraries(kvlang_notify_test PRIVATE kvlang_runtime pthread) +target_include_directories(kvlang_notify_test PRIVATE include src) +add_executable(kvlang_incr_test tests/incr_test.c) +target_link_libraries(kvlang_incr_test PRIVATE kvlang_runtime pthread) +target_include_directories(kvlang_incr_test PRIVATE include src) +add_executable(kvlang_expire_test tests/expire_test.c) +target_link_libraries(kvlang_expire_test PRIVATE kvlang_runtime pthread) +target_include_directories(kvlang_expire_test PRIVATE include src) +add_executable(kvlang_watchany_test tests/watchany_test.c) +target_link_libraries(kvlang_watchany_test PRIVATE kvlang_runtime pthread) +target_include_directories(kvlang_watchany_test PRIVATE include src) +if(KVSPACE_LIB STREQUAL "kvspace-c") + set_target_properties(kvlang_delegate_test PROPERTIES + BUILD_RPATH "${KVSPACE_C_DIR};${BLOCKMALLOC_DIR};${SLOTSBOXMALLOC_DIR};${BIN_DIR}") + set_target_properties(kvlang_notify_test PROPERTIES + BUILD_RPATH "${KVSPACE_C_DIR};${BLOCKMALLOC_DIR};${SLOTSBOXMALLOC_DIR};${BIN_DIR}") + set_target_properties(kvlang_incr_test PROPERTIES + BUILD_RPATH "${KVSPACE_C_DIR};${BLOCKMALLOC_DIR};${SLOTSBOXMALLOC_DIR};${BIN_DIR}") + set_target_properties(kvlang_expire_test PROPERTIES + BUILD_RPATH "${KVSPACE_C_DIR};${BLOCKMALLOC_DIR};${SLOTSBOXMALLOC_DIR};${BIN_DIR}") + set_target_properties(kvlang_watchany_test PROPERTIES + BUILD_RPATH "${KVSPACE_C_DIR};${BLOCKMALLOC_DIR};${SLOTSBOXMALLOC_DIR};${BIN_DIR}") +else() + set_target_properties(kvlang_delegate_test PROPERTIES BUILD_RPATH "${KVSPACE_DURABLE_DIR};${BIN_DIR}") + set_target_properties(kvlang_notify_test PROPERTIES BUILD_RPATH "${KVSPACE_DURABLE_DIR};${BIN_DIR}") + set_target_properties(kvlang_incr_test PROPERTIES BUILD_RPATH "${KVSPACE_DURABLE_DIR};${BIN_DIR}") + set_target_properties(kvlang_expire_test PROPERTIES BUILD_RPATH "${KVSPACE_DURABLE_DIR};${BIN_DIR}") + set_target_properties(kvlang_watchany_test PROPERTIES BUILD_RPATH "${KVSPACE_DURABLE_DIR};${BIN_DIR}") +endif() diff --git a/runtime/src/dispatch.c b/runtime/src/dispatch.c new file mode 100644 index 00000000..eb712935 --- /dev/null +++ b/runtime/src/dispatch.c @@ -0,0 +1,498 @@ +#include "runtime_internal.h" +#include +#include +#include + +#define HEARTBEAT_STALE_S 15 +#define NEG_CACHE_TTL_NS 250000000LL +#define NEG_CACHE_CAP 64 + +static int64_t default_timeout_ns = 30LL * 1000000000LL; +static int64_t task_status_ttl_ns = 10LL * 60LL * 1000000000LL; + +void kvlangDispatchSetDefaultTimeoutNs(int64_t ns) { default_timeout_ns = ns; } +void kvlangDispatchSetTaskStatusTtlNs(int64_t ns) { task_status_ttl_ns = ns; } + +static int64_t now_ns(void) { + struct timespec t; + clock_gettime(CLOCK_REALTIME, &t); + return (int64_t)t.tv_sec * 1000000000LL + t.tv_nsec; +} + +static char *get_char(kvlangKv_t *kv, const char *key) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangKvGetOne(kv, key, &v); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +static void trim_inplace(char *s) { + if (!s) return; + size_t n = strlen(s); + while (n > 0 && (s[n - 1] == ' ' || s[n - 1] == '\n' || s[n - 1] == '\r' || s[n - 1] == '\t')) s[--n] = 0; + char *p = s; + while (*p == ' ' || *p == '\n' || *p == '\r' || *p == '\t') p++; + if (p != s) memmove(s, p, strlen(p) + 1); +} + +static bool heartbeat_fresh(kvlangKv_t *kv, const char *backend) { + char *k = kvlangKeytreeSysRwirBackendHeartbeat(backend); + char *s = get_char(kv, k); + free(k); + if (!s) return false; + trim_inplace(s); + char *end = NULL; + long long sec = strtoll(s, &end, 10); + bool bad = end == s || (end && *end); + free(s); + if (bad) return false; + time_t now = time(NULL); + long long d = (long long)now - sec; + if (d < 0) d = -d; + return d <= HEARTBEAT_STALE_S; +} + +static bool is_on_duty(kvlangKv_t *kv, const char *backend) { + char *k = kvlangKeytreeSysRwirBackendStatus(backend); + char *s = get_char(kv, k); + free(k); + if (!s) return false; + trim_inplace(s); + bool ok = (strcmp(s, "ready") == 0 || strcmp(s, "busy") == 0) && heartbeat_fresh(kv, backend); + free(s); + return ok; +} + +static double parse_load(kvlangKv_t *kv, const char *backend) { + char *k = kvlangKeytreeSysRwirBackendLoad(backend); + char *s = get_char(kv, k); + free(k); + if (!s) return 0; + trim_inplace(s); + char *end = NULL; + double f = strtod(s, &end); + bool bad = end == s || (end && *end) || isnan(f) || isinf(f) || f < 0; + free(s); + return bad ? 1.0 : f; +} + +typedef struct { char **v; int n, cap; } strlist_t; + +static void sl_push(strlist_t *l, char *s) { + if (l->n == l->cap) { + l->cap = l->cap ? l->cap * 2 : 8; + l->v = realloc(l->v, sizeof(char *) * (size_t)l->cap); + } + l->v[l->n++] = s; +} + +static void sl_free(strlist_t *l) { + for (int i = 0; i < l->n; i++) free(l->v[i]); + free(l->v); + l->v = NULL; l->n = l->cap = 0; +} + +static void list_children(kvlangKv_t *kv, const char *dir, strlist_t *out) { + char *pref = kvlangKeytreeStack(dir); + char **names = NULL; int n = 0; + kvlangKvList(kv, pref, false, false, &names, &n); + free(pref); + for (int i = 0; i < n; i++) { + size_t ln = strlen(names[i]); + if (ln > 0 && names[i][ln - 1] == '/') names[i][ln - 1] = 0; + sl_push(out, names[i]); + } + free(names); +} + +static void backends_for(kvlangKv_t *kv, const char *op, strlist_t *reg, strlist_t *duty) { + if (!kvlangKeytreeValidSegment(op)) return; + char *root = kvlangKeytreeSysRwirBackendRoot(); + strlist_t names = {0}; + list_children(kv, root, &names); + free(root); + for (int i = 0; i < names.n; i++) { + char *opk = kvlangKeytreeSysRwirBackendOp(names.v[i], op); + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangKvGetOne(kv, opk, &v); + free(opk); + if (kvlangXvalueNone(&v)) { kvlangXvalueFree(&v); continue; } + kvlangXvalueFree(&v); + sl_push(reg, strdup(names.v[i])); + if (is_on_duty(kv, names.v[i])) sl_push(duty, strdup(names.v[i])); + } + sl_free(&names); +} + +typedef struct { kvlangKv_t *kv; char op[96]; int64_t exp; int used; } negent_t; + +static struct { + pthread_mutex_t mu; + negent_t e[NEG_CACHE_CAP]; +} neg = { .mu = PTHREAD_MUTEX_INITIALIZER }; + +static bool neg_hit(kvlangKv_t *kv, const char *op) { + int64_t t = now_ns(); + pthread_mutex_lock(&neg.mu); + bool hit = false; + for (int i = 0; i < NEG_CACHE_CAP; i++) { + if (!neg.e[i].used) continue; + if (neg.e[i].kv != kv || strcmp(neg.e[i].op, op) != 0) continue; + if (t >= neg.e[i].exp) { neg.e[i].used = 0; break; } + hit = true; + break; + } + pthread_mutex_unlock(&neg.mu); + return hit; +} + +static void neg_put(kvlangKv_t *kv, const char *op) { + int64_t exp = now_ns() + NEG_CACHE_TTL_NS; + pthread_mutex_lock(&neg.mu); + int slot = 0; + for (int i = 0; i < NEG_CACHE_CAP; i++) { + if (!neg.e[i].used || (neg.e[i].kv == kv && strcmp(neg.e[i].op, op) == 0)) { slot = i; break; } + } + neg.e[slot].kv = kv; + snprintf(neg.e[slot].op, sizeof neg.e[slot].op, "%s", op); + neg.e[slot].exp = exp; + neg.e[slot].used = 1; + pthread_mutex_unlock(&neg.mu); +} + +bool kvlangDispatchIsDelegatedOp(kvlangKv_t *kv, const char *opcode) { + char *op = kvlangKeytreeCanonOp(opcode); + if (neg_hit(kv, op)) { free(op); return false; } + strlist_t reg = {0}, duty = {0}; + backends_for(kv, op, ®, &duty); + bool ok = duty.n > 0; + if (!ok) neg_put(kv, op); + sl_free(®); sl_free(&duty); free(op); + return ok; +} + +static int select_backend(kvlangKv_t *kv, const char *op, char **out, char *err, size_t err_cap) { + *out = NULL; + strlist_t reg = {0}, duty = {0}; + backends_for(kv, op, ®, &duty); + if (reg.n == 0) { + snprintf(err, err_cap, "no backend supports opcode=%s", op); + sl_free(®); sl_free(&duty); + return -1; + } + if (duty.n == 0) { + kvlangStrbuf_t got; kvlangStrbufInit(&got); + for (int i = 0; i < reg.n; i++) { + if (i) kvlangStrbufPuts(&got, ", "); + char *sk = kvlangKeytreeSysRwirBackendStatus(reg.v[i]); + char *st = get_char(kv, sk); + free(sk); + kvlangStrbufPuts(&got, reg.v[i]); + kvlangStrbufPuts(&got, "="); + kvlangStrbufPuts(&got, st ? st : ""); + free(st); + } + snprintf(err, err_cap, "no on-duty backend for opcode=%s; backend status must be ready|busy, got %s", op, got.p ? got.p : ""); + kvlangStrbufFree(&got); + sl_free(®); sl_free(&duty); + return -1; + } + char *best = duty.v[0]; + double best_load = parse_load(kv, best); + for (int i = 1; i < duty.n; i++) { + double ld = parse_load(kv, duty.v[i]); + if (ld < best_load) { best_load = ld; best = duty.v[i]; } + } + *out = strdup(best); + sl_free(®); sl_free(&duty); + return 0; +} + +static int64_t timeout_for(kvlangKv_t *kv, const char *backend) { + int64_t d = default_timeout_ns; + char *root = kvlangKeytreeSysRwirBackendCategoryRoot(backend); + strlist_t cats = {0}; + list_children(kv, root, &cats); + free(root); + for (int i = 0; i < cats.n; i++) { + int64_t t = 0; + if (strcmp(cats.v[i], "compute") == 0) t = 30LL * 1000000000LL; + else if (strcmp(cats.v[i], "api") == 0) t = 120LL * 1000000000LL; + else if (strcmp(cats.v[i], "agent") == 0) t = 300LL * 1000000000LL; + if (t > d) d = t; + } + sl_free(&cats); + return d; +} + +static int set_char(kvlangKv_t *kv, const char *key, const char *s, char *err, uint32_t err_cap) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewCharUtf8(&v, s); + kvlangKvPair_t p = { (char *)key, v }; + int rc = kvlangKvSet(kv, &p, 1, err, err_cap); + kvlangXvalueFree(&v); + return rc; +} + +static void json_esc(kvlangStrbuf_t *b, const char *s) { + kvlangStrbufPutc(b, '"'); + for (; s && *s; s++) { + unsigned char c = (unsigned char)*s; + if (c == '"' || c == '\\') { kvlangStrbufPutc(b, '\\'); kvlangStrbufPutc(b, (char)c); } + else if (c == '\n') kvlangStrbufPuts(b, "\\n"); + else if (c == '\r') kvlangStrbufPuts(b, "\\r"); + else if (c == '\t') kvlangStrbufPuts(b, "\\t"); + else kvlangStrbufPutc(b, (char)c); + } + kvlangStrbufPutc(b, '"'); +} + +static bool is_var_slot(const kvlangParam_t *p) { + if (!p->name || !p->name[0]) return false; + if (p->name[0] == '/') return true; + return kvlangXvalueKindIs(&p->val, KVSPACE_KIND_RWIR) || kvlangXvalueNone(&p->val); +} + +typedef struct { char **keys; int n; } keyset_t; + +static bool keyset_has(const keyset_t *s, const char *k) { + for (int i = 0; i < s->n; i++) if (strcmp(s->keys[i], k) == 0) return true; + return false; +} + +static int fail_msg(kvlangKv_t *kv, const char *vtid, const char *pc, const char *fmt, ...) { + char body[384]; + va_list ap; va_start(ap, fmt); + vsnprintf(body, sizeof body, fmt, ap); + va_end(ap); + char msg[420]; + snprintf(msg, sizeof msg, "RuntimeError: delegate: %s", body); + kvlangVthreadSetError(kv, vtid, pc, msg); + return KVLANG_DELEGATE_ERR; +} + +static int fail_as(kvlangKv_t *kv, const char *vtid, const char *pc, const char *msg) { + kvlangVthreadSetError(kv, vtid, pc, msg); + return KVLANG_DELEGATE_ERR; +} + +static int lib_arity(const kvlangXvalue_t *v, int *nr, int *nw) { + kvspaceHead_t h; + if (kvlangXvalueHead(v, &h) != 0) return -1; + int32_t bl = 0; + const uint8_t *b = kvlangXvalueBody(v, &h, &bl); + if (!b || bl < 4) return -1; + *nr = b[0] | (b[1] << 8); + *nw = b[2] | (b[3] << 8); + return 0; +} + +int kvlangDispatchDelegate(kvlangKv_t *kv, const char *vtid, const char *pc, kvlangRwirInst_t *inst) { + char *op = kvlangKeytreeCanonOp(inst->opcode); + + char *declk = kvlangKeytreeRwir(op); + kvlangXvalue_t decl; kvlangXvalueZero(&decl); + kvlangKvGetOne(kv, declk, &decl); + free(declk); + + char *sigk = kvlangKeytreeLibSig(op); + kvlangXvalue_t sig; kvlangXvalueZero(&sig); + kvlangKvGetOne(kv, sigk, &sig); + free(sigk); + + if (!kvlangXvalueNone(&decl) && kvlangXvalueKindIs(&decl, KVSPACE_KIND_RWIR)) { + int dnr = 0, dnw = 0; + if (lib_arity(&decl, &dnr, &dnw) == 0) { + if (inst->nr != dnr) { + kvlangXvalueFree(&decl); kvlangXvalueFree(&sig); + int rc = fail_msg(kv, vtid, pc, "%s: read arity mismatch: call site has %d, declaration has %d", op, inst->nr, dnr); + free(op); + return rc; + } + if (inst->nw != dnw) { + kvlangXvalueFree(&decl); kvlangXvalueFree(&sig); + int rc = fail_msg(kv, vtid, pc, "%s: write arity mismatch: call site has %d, declaration has %d", op, inst->nw, dnw); + free(op); + return rc; + } + } + } else if (!kvlangXvalueNone(&sig) && kvlangXvalueKindIs(&sig, KVSPACE_KIND_RWFUNC)) { + strlist_t reg = {0}, duty = {0}; + backends_for(kv, op, ®, &duty); + sl_free(®); + if (duty.n == 0) { + sl_free(&duty); + kvlangXvalueFree(&decl); kvlangXvalueFree(&sig); + free(op); + return KVLANG_DELEGATE_LOCAL; + } + kvlangStrbuf_t names; kvlangStrbufInit(&names); + for (int i = 0; i < duty.n; i++) { + if (i) kvlangStrbufPuts(&names, ","); + kvlangStrbufPuts(&names, duty.v[i]); + } + sl_free(&duty); + kvlangXvalueFree(&decl); kvlangXvalueFree(&sig); + int rc = fail_msg(kv, vtid, pc, "%s: on-duty backend %s and local rwfunc both define it; one opcode, one definition", op, names.p ? names.p : ""); + kvlangStrbufFree(&names); + free(op); + return rc; + } else if (!kvlangXvalueNone(&decl) && !kvlangXvalueKindIs(&decl, KVSPACE_KIND_RWIR)) { + const char *k = kvlangXvalueKind(&decl); + int rc = fail_msg(kv, vtid, pc, "%s: /lib entry holds kind %s, which is not callable", op, k[0] ? k : "?"); + kvlangXvalueFree(&decl); kvlangXvalueFree(&sig); + free(op); + return rc; + } + kvlangXvalueFree(&decl); + kvlangXvalueFree(&sig); + + char selerr[256]; + char *backend = NULL; + if (select_backend(kv, op, &backend, selerr, sizeof selerr) != 0) { + int rc = fail_msg(kv, vtid, pc, "%s", selerr); + free(op); + return rc; + } + + char *seqk = kvlangKeytreeVthreadDelegSeq(vtid); + int64_t seq = kvlangVthreadNextSeq(kv, seqk); + free(seqk); + char task_id[160]; + snprintf(task_id, sizeof task_id, "rwir:%s:%s:%lld", backend, vtid, (long long)seq); + + char *fr = kvlangKeytreeFrameRoot(pc); + char **out_keys = calloc((size_t)inst->nw, sizeof(char *)); + keyset_t outset = { out_keys, 0 }; + for (int i = 0; i < inst->nw; i++) { + char *key = kvlangBuiltinResolveWriteSlot(kv, fr, inst->writes[i].name); + char *werr = kvlangKeytreeCheckWriteKey(vtid, key); + if (werr) { + free(key); + for (int j = 0; j < i; j++) free(out_keys[j]); + free(out_keys); free(fr); free(backend); free(op); + int rc = fail_as(kv, vtid, pc, werr); + free(werr); + return rc; + } + out_keys[i] = key; + outset.n++; + } + + kvlangStrbuf_t json; kvlangStrbufInit(&json); + kvlangStrbufPuts(&json, "{\"request_id\":"); json_esc(&json, task_id); + kvlangStrbufPuts(&json, ",\"vtid\":"); json_esc(&json, vtid); + kvlangStrbufPuts(&json, ",\"pc\":"); json_esc(&json, pc); + kvlangStrbufPuts(&json, ",\"opcode\":"); json_esc(&json, op); + kvlangStrbufPuts(&json, ",\"inputs\":["); + for (int i = 0; i < inst->nr; i++) { + if (i) kvlangStrbufPutc(&json, ','); + kvlangStrbufPutc(&json, '{'); + if (is_var_slot(&inst->reads[i]) && inst->reads[i].name[0] == '/') { + kvlangStrbufPuts(&json, "\"key\":"); json_esc(&json, inst->reads[i].name); + } else if (is_var_slot(&inst->reads[i])) { + char *rk = kvlangBuiltinResolveWriteSlot(kv, fr, inst->reads[i].name); + if (keyset_has(&outset, rk)) { + kvlangXvalue_t val; kvlangXvalueZero(&val); + kvlangBuiltinResolveReadValue(kv, fr, inst->reads[i].name, &inst->reads[i].val, &val); + char *disp = NULL; kvlangDisplay(&val, &disp); + kvlangStrbufPuts(&json, "\"value\":"); json_esc(&json, disp ? disp : ""); + free(disp); kvlangXvalueFree(&val); + } else { + kvlangStrbufPuts(&json, "\"key\":"); json_esc(&json, rk); + } + free(rk); + } else { + char *disp = NULL; kvlangDisplay(&inst->reads[i].val, &disp); + kvlangStrbufPuts(&json, "\"value\":"); json_esc(&json, disp ? disp : ""); + free(disp); + } + kvlangStrbufPutc(&json, '}'); + } + kvlangStrbufPuts(&json, "],\"outputs\":["); + for (int i = 0; i < inst->nw; i++) { + if (i) kvlangStrbufPutc(&json, ','); + kvlangStrbufPuts(&json, "{\"key\":"); json_esc(&json, out_keys[i]); + kvlangStrbufPutc(&json, '}'); + } + char *done_key = kvlangKeytreeDoneRwir(task_id); + kvlangStrbufPuts(&json, "],\"done_key\":"); json_esc(&json, done_key); + kvlangStrbufPutc(&json, '}'); + + char *status_key = kvlangKeytreeSysTask(task_id, SEG_STATUS); + char err[256]; + if (set_char(kv, status_key, "pending", err, sizeof err) != 0) { + fail_msg(kv, vtid, pc, "%s: init task status: %s", op, err); + goto done_err; + } + + for (int i = 0; i < inst->nw; i++) { + if (kvlangKvDel(kv, out_keys[i], err, sizeof err) != 0) { + fail_msg(kv, vtid, pc, "%s: clear output slot %s: %s", op, out_keys[i], err); + goto done_err; + } + } + + char *cmd = kvlangKeytreeSysRwirBackendCmd(backend); + kvlangXvalue_t cmdv; kvlangXvalueZero(&cmdv); + kvlangXvalueNewCharUtf8(&cmdv, json.p); + if (kvlangKvNotify(kv, cmd, &cmdv, err, sizeof err) != 0) { + kvlangXvalueFree(&cmdv); free(cmd); + fail_msg(kv, vtid, pc, "push task: %s", err); + goto done_err; + } + kvlangXvalueFree(&cmdv); + free(cmd); + + kvlangVthreadSet(kv, vtid, pc, "wait"); + kvlangLogDebug("[%s] DELEGATE %s task=%s", vtid, inst->opcode, task_id); + + kvlangXvalue_t want; kvlangXvalueZero(&want); + kvlangXvalueNewCharUtf8(&want, "done"); + kvlangXvalue_t got; kvlangXvalueZero(&got); + kvlangKvWatch(kv, status_key, &want, (uint64_t)timeout_for(kv, backend), &got); + kvlangXvalueFree(&want); kvlangXvalueFree(&got); + + char *st = get_char(kv, status_key); + if (!st || strcmp(st, "done") != 0) { + fail_msg(kv, vtid, pc, "%s: timeout or failed (status=%s)", op, st ? st : ""); + free(st); + goto done_err; + } + free(st); + + for (int i = 0; i < inst->nw; i++) { + kvlangXvalue_t ov; kvlangXvalueZero(&ov); + kvlangKvGetOne(kv, out_keys[i], &ov); + bool empty = kvlangXvalueNone(&ov); + kvlangXvalueFree(&ov); + if (empty) { + fail_msg(kv, vtid, pc, "%s: executor reported done but left output slot %s unwritten", op, out_keys[i]); + goto done_err; + } + } + + kvlangStrbuf_t npc; kvlangStrbufInit(&npc); + kvlangRwirNextPc(pc, &npc); + kvlangVthreadSet(kv, vtid, npc.p, "running"); + kvlangStrbufFree(&npc); + kvlangLogDebug("[%s] DONE %s task=%s", vtid, inst->opcode, task_id); + + if (kvlangKvExpire(kv, status_key, (uint64_t)task_status_ttl_ns, err, sizeof err) != 0) { + fail_msg(kv, vtid, pc, "%s: expire task status: %s", op, err); + goto done_err; + } + + for (int i = 0; i < inst->nw; i++) free(out_keys[i]); + free(out_keys); free(fr); free(backend); free(op); + free(done_key); free(status_key); kvlangStrbufFree(&json); + return KVLANG_DELEGATE_OK; + +done_err: + if (status_key) kvlangKvExpire(kv, status_key, (uint64_t)task_status_ttl_ns, err, sizeof err); + for (int i = 0; i < inst->nw; i++) free(out_keys[i]); + free(out_keys); free(fr); free(backend); free(op); + free(done_key); free(status_key); kvlangStrbufFree(&json); + return KVLANG_DELEGATE_ERR; +} diff --git a/runtime/src/keytree.c b/runtime/src/keytree.c index 7a0a5977..7a3ee1a3 100644 --- a/runtime/src/keytree.c +++ b/runtime/src/keytree.c @@ -127,3 +127,153 @@ bool kvlangKeytreeIsEntryPc(const char *pc) { const char *slash = strrchr(pc, '/'); return slash && strcmp(slash, "/[1,0]") == 0; } + +char *kvlangKeytreeCanonOp(const char *opcode) { + const char *p = opcode ? opcode : ""; + if (strncmp(p, LIB_ROOT PATH_SEP, 5) == 0) p += 5; + return strdup(p); +} + +bool kvlangKeytreeValidSegment(const char *s) { + if (!s || !s[0]) return false; + if (strcmp(s, ".") == 0 || strcmp(s, "..") == 0) return false; + if (strcmp(s, MEMBER_SEP) == 0) return false; + if (s[0] == '.' && s[1] == '.' && s[2] == 0) return false; + return strchr(s, '/') == NULL; +} + +static char *join3(const char *a, const char *b, const char *c) { + kvlangStrbuf_t s; kvlangStrbufInit(&s); + kvlangStrbufPuts(&s, a); + if (b && b[0]) { kvlangStrbufPutc(&s, '/'); kvlangStrbufPuts(&s, b); } + if (c && c[0]) { kvlangStrbufPutc(&s, '/'); kvlangStrbufPuts(&s, c); } + return kvlangStrbufDetach(&s); +} + +char *kvlangKeytreeSysRwirBackendRoot(void) { + return strdup(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND); +} + +char *kvlangKeytreeSysRwirBackend(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, NULL); +} + +char *kvlangKeytreeSysRwirBackendOp(const char *name, const char *opcode) { + kvlangStrbuf_t s; kvlangStrbufInit(&s); + kvlangStrbufPuts(&s, SYS_ROOT PATH_SEP SEG_RWIR_BACKEND "/"); + kvlangStrbufPuts(&s, name); + kvlangStrbufPuts(&s, "/" SEG_OP "/"); + kvlangStrbufPuts(&s, opcode); + return kvlangStrbufDetach(&s); +} + +char *kvlangKeytreeSysRwirBackendCmd(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, SEG_CMD); +} + +char *kvlangKeytreeSysRwirBackendStatus(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, SEG_STATUS); +} + +char *kvlangKeytreeSysRwirBackendLoad(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, SEG_LOAD); +} + +char *kvlangKeytreeSysRwirBackendHeartbeat(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, SEG_LAST_HEARTBEAT); +} + +char *kvlangKeytreeSysRwirBackendCategoryRoot(const char *name) { + return join3(SYS_ROOT PATH_SEP SEG_RWIR_BACKEND, name, SEG_CATEGORY); +} + +char *kvlangKeytreeSysTask(const char *task_id, const char *field) { + kvlangStrbuf_t s; kvlangStrbufInit(&s); + kvlangStrbufPuts(&s, SYS_ROOT PATH_SEP SEG_TASK "/"); + kvlangStrbufPuts(&s, task_id); + char *base = kvlangStrbufDetach(&s); + char *r = kvlangKeytreeMember(base, field); + free(base); + return r; +} + +char *kvlangKeytreeDoneRwir(const char *task_id) { + return join3(DONE_ROOT PATH_SEP SEG_RWIR, task_id, NULL); +} + +char *kvlangKeytreeVthreadDelegSeq(const char *vtid) { + kvlangStrbuf_t s; kvlangStrbufInit(&s); + kvlangKeytreeVthread(vtid, &s); + kvlangStrbufPutc(&s, '/'); + kvlangStrbufPuts(&s, RUNTIME_MEMBER_SEP); + kvlangStrbufPuts(&s, SEG_DELEGSEQ); + return kvlangStrbufDetach(&s); +} + +char *kvlangKeytreeLibSig(const char *opcode) { + kvlangStrbuf_t s; kvlangStrbufInit(&s); + kvlangStrbufPuts(&s, LIB_ROOT PATH_SEP); + kvlangStrbufPuts(&s, opcode); + kvlangStrbufPuts(&s, "/[0,0]"); + return kvlangStrbufDetach(&s); +} + +char *kvlangKeytreeCheckWriteKey(const char *vtid, const char *key) { + if (!key || key[0] != '/') { + char *e = malloc(160); + snprintf(e, 160, "ValueError: write slot is not an absolute path: \"%s\"; help: a write slot must resolve to a canonical path starting with /", key ? key : ""); + return e; + } + const char *p = key + 1; + while (*p) { + const char *slash = strchr(p, '/'); + size_t n = slash ? (size_t)(slash - p) : strlen(p); + if (n == 0 || (n == 1 && p[0] == '.') || (n == 2 && p[0] == '.' && p[1] == '.')) { + char *e = malloc(160); + snprintf(e, 160, "ValueError: write slot is not a canonical path: \"%s\"; help: a path may not contain an empty segment, . or ..", key); + return e; + } + if (n >= 3 && (unsigned char)p[0] == 0xE2 && (unsigned char)p[1] == 0x80 && (unsigned char)p[2] == 0xA5) { + char *e = malloc(192); + snprintf(e, 192, "PermissionError: write slot targets an engine-reserved key: \"%s\"; help: keys starting with %s belong to the VM and programs cannot write them", key, RUNTIME_MEMBER_SEP); + return e; + } + p = slash ? slash + 1 : p + n; + } + static const char *roots[] = { LIB_ROOT, SYS_ROOT, DEV_ROOT, DONE_ROOT, VTHREAD_ROOT, NULL }; + for (int i = 0; roots[i]; i++) { + if (strcmp(key, roots[i]) == 0) { + char *e = malloc(192); + snprintf(e, 192, "PermissionError: write slot targets domain root %s: \"%s\"; help: the domain belongs to the VM, and a root is a directory besides", roots[i], key); + return e; + } + } + static const char *prot[] = { LIB_ROOT, SYS_ROOT, DEV_ROOT, DONE_ROOT, NULL }; + for (int i = 0; prot[i]; i++) { + size_t pl = strlen(prot[i]); + if (strncmp(key, prot[i], pl) == 0 && key[pl] == '/') { + char *e = malloc(220); + snprintf(e, 220, "PermissionError: write slot targets protected domain %s: \"%s\"; help: %s belongs to the VM; write a slot in your own vthread or a user global key", prot[i], key, prot[i]); + return e; + } + } + if (strncmp(key, VTHREAD_ROOT "/", strlen(VTHREAD_ROOT) + 1) == 0) { + kvlangStrbuf_t own; kvlangStrbufInit(&own); + kvlangKeytreeVthread(vtid, &own); + if (strcmp(key, own.p) == 0) { + kvlangStrbufFree(&own); + char *e = malloc(200); + snprintf(e, 200, "IsADirectoryError: write slot targets this vthread's frame root: \"%s\"; help: a frame root is a directory and writing it as a leaf destroys the call frame", key); + return e; + } + size_t ol = strlen(own.p); + if (!(strncmp(key, own.p, ol) == 0 && key[ol] == '/')) { + char *e = malloc(240); + snprintf(e, 240, "PermissionError: write slot is outside this vthread's subtree (vtid=%s): \"%s\"; help: the whole /vthread domain belongs to the engine, a program may only write keys under %s/", vtid, key, own.p); + kvlangStrbufFree(&own); + return e; + } + kvlangStrbufFree(&own); + } + return NULL; +} diff --git a/runtime/src/kv.c b/runtime/src/kv.c index 0cc94214..0f4be3fc 100644 --- a/runtime/src/kv.c +++ b/runtime/src/kv.c @@ -126,3 +126,42 @@ int kvlangKvWatch(kvlangKv_t *k, const char *key, const kvlangXvalue_t *target, kvspaceBytesFree(d, len); return 0; } + +int kvlangKvNotify(kvlangKv_t *k, const char *key, const kvlangXvalue_t *val, char *err, uint32_t err_cap) { + const uint8_t *p = val->data ? val->data : (const uint8_t *)""; + return kvspaceNotify(k->h, key, p, val->len, err, err_cap); +} + +int kvlangKvTake(kvlangKv_t *k, const char *key, uint64_t timeout_ns, kvlangXvalue_t *out) { + kvlangXvalueZero(out); + uint8_t *d; uint32_t len; + if (kvspaceTake(k->h, key, timeout_ns, &d, &len) != 0) return -1; + kvlangXvalueCopyMalloc(out, d, len); + kvspaceBytesFree(d, len); + return 0; +} + +int kvlangKvIncr(kvlangKv_t *k, const char *key, int64_t *out, char *err, uint32_t err_cap) { + return kvspaceIncr(k->h, key, out, err, err_cap); +} + +int kvlangKvExpire(kvlangKv_t *k, const char *key, uint64_t ttl_ns, char *err, uint32_t err_cap) { + return kvspaceExpire(k->h, key, ttl_ns, err, err_cap); +} + +int kvlangKvWatchAny(kvlangKv_t *k, const char *const *keys, int n, uint64_t timeout_ns, + char **out_key, kvlangXvalue_t *out) { + *out_key = NULL; + kvlangXvalueZero(out); + uint8_t *kb = NULL, *d = NULL; uint32_t kl = 0, len = 0; + if (kvspaceWatchAny(k->h, keys, (uint32_t)n, timeout_ns, &kb, &kl, &d, &len) != 0) return -1; + if (kb) { + *out_key = malloc((size_t)kl + 1); + memcpy(*out_key, kb, kl); + (*out_key)[kl] = 0; + kvspaceBytesFree(kb, kl); + } + kvlangXvalueCopyMalloc(out, d, len); + kvspaceBytesFree(d, len); + return 0; +} diff --git a/runtime/src/kvcpu.c b/runtime/src/kvcpu.c index 1ec1568b..9d6d5602 100644 --- a/runtime/src/kvcpu.c +++ b/runtime/src/kvcpu.c @@ -596,6 +596,21 @@ int kvlangKvcpuExecuteMode(kvlangKv_t *kv, const char *pc, kvmode_t mode, char * exec_err = kvlangBuiltinNative(&f); } else if (is_copy_op(inst.opcode)) { exec_err = kvlangBuiltinExecuteCopy(kv, vtid, cur, &inst); + } else if (kvlangDispatchIsDelegatedOp(kv, inst.opcode)) { + exec_err = kvlangDispatchDelegate(kv, vtid, cur, &inst); + if (exec_err == KVLANG_DELEGATE_LOCAL) { + kvlangRwirInst_t ci; + ci.opcode = strdup(OP_CALL); + ci.nr = inst.nr + 1; + ci.nw = inst.nw; + ci.reads = malloc(sizeof(kvlangParam_t) * (size_t)ci.nr); + ci.reads[0].name = strdup(inst.opcode); + ci.reads[0].val.data = NULL; ci.reads[0].val.len = 0; + for (int i = 0; i < inst.nr; i++) ci.reads[i + 1] = inst.reads[i]; + ci.writes = inst.writes; + exec_err = handle_control(kv, vtid, cur, &ci); + free(ci.opcode); free(ci.reads[0].name); free(ci.reads); + } } else if (is_ext_rwir(kv, inst.opcode)) { char *rk = kvlangKeytreeRwir(inst.opcode); int def_nr = 0; diff --git a/runtime/src/runtime_internal.h b/runtime/src/runtime_internal.h index 2e1ad447..aac68114 100644 --- a/runtime/src/runtime_internal.h +++ b/runtime/src/runtime_internal.h @@ -36,6 +36,14 @@ extern int kvspaceMkindexExt(void *h, const char *path, const char *ext_path, extern int kvspaceRmindexExt(void *h, const char *path, char *err, uint32_t err_cap); extern int kvspaceWatch(void *h, const char *key, const uint8_t *target, uint32_t target_len, uint64_t tick_ns, uint8_t **out, uint32_t *out_len); +extern int kvspaceNotify(void *h, const char *key, const uint8_t *val, uint32_t len, + char *err, uint32_t err_cap); +extern int kvspaceTake(void *h, const char *key, uint64_t timeout_ns, + uint8_t **out, uint32_t *out_len); +extern int kvspaceIncr(void *h, const char *key, int64_t *out, char *err, uint32_t err_cap); +extern int kvspaceExpire(void *h, const char *key, uint64_t ttl_ns, char *err, uint32_t err_cap); +extern int kvspaceWatchAny(void *h, const char *const *keys, uint32_t nkeys, uint64_t timeout_ns, + uint8_t **out_key, uint32_t *out_key_len, uint8_t **out, uint32_t *out_len); extern int kvspaceTlvEncode(const char *kind, const uint8_t *raw, uint32_t raw_len, const int32_t *dims, int32_t ndim, uint8_t **out, uint32_t *out_len); extern int kvspaceTlvEncodePtr(const char *kind, const uint8_t *raw, uint32_t raw_len, @@ -165,6 +173,12 @@ int kvlangKvDelExtIndex(kvlangKv_t *k, const char *path, char *err, uint32_t err int kvlangKvList(kvlangKv_t *k, const char *prefix, bool expand_ext, bool resolve, char ***out_names, int *out_count); /* split \n */ int kvlangKvWatch(kvlangKv_t *k, const char *key, const kvlangXvalue_t *target, uint64_t tick_ns, kvlangXvalue_t *out); +int kvlangKvNotify(kvlangKv_t *k, const char *key, const kvlangXvalue_t *val, char *err, uint32_t err_cap); +int kvlangKvTake(kvlangKv_t *k, const char *key, uint64_t timeout_ns, kvlangXvalue_t *out); +int kvlangKvIncr(kvlangKv_t *k, const char *key, int64_t *out, char *err, uint32_t err_cap); +int kvlangKvExpire(kvlangKv_t *k, const char *key, uint64_t ttl_ns, char *err, uint32_t err_cap); +int kvlangKvWatchAny(kvlangKv_t *k, const char *const *keys, int n, uint64_t timeout_ns, + char **out_key, kvlangXvalue_t *out); /* ── keytree ───────────────────────────────────────────────────────── */ @@ -177,6 +191,19 @@ int kvlangKvWatch(kvlangKv_t *k, const char *key, const kvlangXvalue_t *target, #define SEG_MSG "msg" #define LIB_ROOT "/lib" #define VTHREAD_ROOT "/vthread" +#define SYS_ROOT "/sys" +#define DEV_ROOT "/dev" +#define DONE_ROOT "/done" +#define SEG_DELEGSEQ "delegseq" +#define SEG_RWIR_BACKEND "rwir-backend" +#define SEG_LAST_HEARTBEAT "last_heartbeat" +#define SEG_CATEGORY "category" +#define SEG_LOAD "load" +#define SEG_CMD "cmd" +#define SEG_OP "op" +#define SEG_TASK "task" +#define SEG_DONE "done" +#define SEG_RWIR "rwir" static inline void kvlangStrbufClear(kvlangStrbuf_t *b) { b->len = 0; if (b->p) b->p[0] = 0; } @@ -200,6 +227,22 @@ void kvlangKeytreeFrameReturnpc(const char *root, kvlangStrbuf_t *out); void kvlangKeytreeFrameRo(const char *root, kvlangStrbuf_t *out); bool kvlangKeytreeIsEntryPc(const char *pc); +char *kvlangKeytreeCanonOp(const char *opcode); /* malloc */ +bool kvlangKeytreeValidSegment(const char *s); +char *kvlangKeytreeCheckWriteKey(const char *vtid, const char *key); /* malloc err or NULL */ +char *kvlangKeytreeSysRwirBackendRoot(void); /* malloc */ +char *kvlangKeytreeSysRwirBackend(const char *name); /* malloc */ +char *kvlangKeytreeSysRwirBackendOp(const char *name, const char *opcode); +char *kvlangKeytreeSysRwirBackendCmd(const char *name); +char *kvlangKeytreeSysRwirBackendStatus(const char *name); +char *kvlangKeytreeSysRwirBackendLoad(const char *name); +char *kvlangKeytreeSysRwirBackendHeartbeat(const char *name); +char *kvlangKeytreeSysRwirBackendCategoryRoot(const char *name); +char *kvlangKeytreeSysTask(const char *task_id, const char *field); +char *kvlangKeytreeDoneRwir(const char *task_id); +char *kvlangKeytreeVthreadDelegSeq(const char *vtid); +char *kvlangKeytreeLibSig(const char *opcode); /* /lib//[0,0] */ + /* ── rwir ──────────────────────────────────────────────────────────── */ #define OP_CALL "call" @@ -232,6 +275,18 @@ void kvlangVthreadGet(kvlangKv_t *kv, const char *vtid, char **pc, char **status void kvlangVthreadSet(kvlangKv_t *kv, const char *vtid, const char *pc, const char *status); void kvlangVthreadSetDone(kvlangKv_t *kv, const char *vtid, const char *ret); void kvlangVthreadSetError(kvlangKv_t *kv, const char *vtid, const char *pc, const char *msg); +int64_t kvlangVthreadNextSeq(kvlangKv_t *kv, const char *key); + +/* ── dispatch (rwir-backend) ───────────────────────────────────────── */ + +#define KVLANG_DELEGATE_OK 0 +#define KVLANG_DELEGATE_ERR -1 +#define KVLANG_DELEGATE_LOCAL 1 /* backend left duty; run local rwfunc */ + +bool kvlangDispatchIsDelegatedOp(kvlangKv_t *kv, const char *opcode); +int kvlangDispatchDelegate(kvlangKv_t *kv, const char *vtid, const char *pc, kvlangRwirInst_t *inst); +void kvlangDispatchSetDefaultTimeoutNs(int64_t ns); +void kvlangDispatchSetTaskStatusTtlNs(int64_t ns); /* ── builtin ───────────────────────────────────────────────────────── */ diff --git a/runtime/src/vthread.c b/runtime/src/vthread.c index 52c708f8..fb973915 100644 --- a/runtime/src/vthread.c +++ b/runtime/src/vthread.c @@ -67,3 +67,13 @@ void kvlangVthreadSetError(kvlangKv_t *kv, const char *vtid, const char *pc, con kvlangXvalueFree(&vpc); kvlangXvalueFree(&vmsg); kvlangXvalueFree(&vst); kvlangStrbufFree(&msg_path); kvlangStrbufFree(&pc_key); kvlangStrbufFree(&st_key); } + +int64_t kvlangVthreadNextSeq(kvlangKv_t *kv, const char *key) { + int64_t n = 0; + char err[256]; + if (kvlangKvIncr(kv, key, &n, err, sizeof err) != 0) { + fprintf(stderr, "NextSeq: %s\n", err[0] ? err : "Incr failed"); + abort(); + } + return n; +} diff --git a/runtime/tests/delegate_test.c b/runtime/tests/delegate_test.c new file mode 100644 index 00000000..c4056aba --- /dev/null +++ b/runtime/tests/delegate_test.c @@ -0,0 +1,310 @@ +#include "runtime_internal.h" +#include +#include +#include + +static void fail(const char *msg) { + fprintf(stderr, "FAIL %s\n", msg); + exit(1); +} + +static void expect(bool ok, const char *msg) { + if (!ok) fail(msg); +} + +static kvlangKv_t *open_kv(void) { + char dsn[128]; + snprintf(dsn, sizeof dsn, "shm:///tmp/kvlang-delegate-%d", (int)getpid()); + kvlangKv_t *kv = kvlangKvConnect(dsn); + if (!kv) fail("connect"); + char err[256]; + kvlangKvMkindex(kv, "/sys/rwir-backend/", err, sizeof err); + kvlangKvMkindex(kv, "/sys/task/", err, sizeof err); + kvlangKvMkindex(kv, "/vthread/", err, sizeof err); + kvlangKvMkindex(kv, "/data/", err, sizeof err); + return kv; +} + +static void set_char(kvlangKv_t *kv, const char *key, const char *s) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewCharUtf8(&v, s); + kvlangKvPair_t p = { (char *)key, v }; + char err[256]; + if (kvlangKvSet(kv, &p, 1, err, sizeof err) != 0) fail(err[0] ? err : "set"); + kvlangXvalueFree(&v); +} + +static char *get_char(kvlangKv_t *kv, const char *key) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangKvGetOne(kv, key, &v); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +static char *take_char(kvlangKv_t *kv, const char *key, uint64_t ns) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangKvTake(kv, key, ns, &v); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +static void register_backend(kvlangKv_t *kv, const char *name, const char *opcode, const char *status) { + char *dir = kvlangKeytreeSysRwirBackend(name); + char path[256]; + snprintf(path, sizeof path, "%s/", dir); + char err[256]; + kvlangKvMkindex(kv, path, err, sizeof err); + snprintf(path, sizeof path, "%s/op/", dir); + kvlangKvMkindex(kv, path, err, sizeof err); + free(dir); + char *opk = kvlangKeytreeSysRwirBackendOp(name, opcode); + set_char(kv, opk, "1"); + free(opk); + char *sk = kvlangKeytreeSysRwirBackendStatus(name); + set_char(kv, sk, status); + free(sk); + char *hk = kvlangKeytreeSysRwirBackendHeartbeat(name); + char beat[32]; + snprintf(beat, sizeof beat, "%lld", (long long)time(NULL)); + set_char(kv, hk, beat); + free(hk); +} + +static kvlangRwirInst_t make_inst(const char *op, const char *write_key) { + kvlangRwirInst_t inst; + memset(&inst, 0, sizeof inst); + inst.opcode = strdup(op); + inst.nr = 0; + inst.nw = 1; + inst.reads = NULL; + inst.writes = calloc(1, sizeof(kvlangParam_t)); + inst.writes[0].name = strdup(write_key); + kvlangXvalueZero(&inst.writes[0].val); + return inst; +} + +typedef struct { + kvlangKv_t *kv; + char *cmd; + char *out_key; + char *status_prefix; + int write_out; + int delay_ms; + volatile int ran; +} exec_arg_t; + +static void *executor(void *p) { + exec_arg_t *a = p; + char *json = take_char(a->kv, a->cmd, 2000000000ULL); + if (!json) return NULL; + if (a->delay_ms > 0) usleep((useconds_t)a->delay_ms * 1000); + if (a->write_out) set_char(a->kv, a->out_key, "ok"); + const char *rid = strstr(json, "\"request_id\":\""); + if (!rid) { free(json); a->ran = 1; return NULL; } + rid += 14; + const char *end = strchr(rid, '"'); + if (!end) { free(json); a->ran = 1; return NULL; } + char id[160]; + size_t n = (size_t)(end - rid); + if (n >= sizeof id) n = sizeof id - 1; + memcpy(id, rid, n); id[n] = 0; + free(json); + char *sk = kvlangKeytreeSysTask(id, "status"); + set_char(a->kv, sk, "done"); + free(sk); + a->ran = 1; + return NULL; +} + +static void test_check_write_key(void) { + expect(kvlangKeytreeCheckWriteKey("1", "/data/out") == NULL, "user global allowed"); + char *e = kvlangKeytreeCheckWriteKey("1", "/lib/x"); + expect(e && strncmp(e, "PermissionError", 15) == 0, "/lib rejected"); + free(e); + e = kvlangKeytreeCheckWriteKey("1", "/sys/rwir-backend/fake/cmd"); + expect(e && strncmp(e, "PermissionError", 15) == 0, "/sys rejected"); + free(e); + e = kvlangKeytreeCheckWriteKey("1", "/vthread/2/x"); + expect(e && strncmp(e, "PermissionError", 15) == 0, "other vthread rejected"); + free(e); + e = kvlangKeytreeCheckWriteKey("1", "rel"); + expect(e && strncmp(e, "ValueError", 10) == 0, "relative rejected"); + free(e); +} + +static void test_is_delegated(kvlangKv_t *kv) { + expect(!kvlangDispatchIsDelegatedOp(kv, "fake.echo"), "no backend"); + usleep(300000); + register_backend(kv, "fake", "fake.echo", "ready"); + expect(kvlangDispatchIsDelegatedOp(kv, "fake.echo"), "on duty"); + expect(kvlangDispatchIsDelegatedOp(kv, "/lib/fake.echo"), "folded opcode"); + char *sk = kvlangKeytreeSysRwirBackendStatus("fake"); + set_char(kv, sk, "offline"); + free(sk); + usleep(300000); + expect(!kvlangDispatchIsDelegatedOp(kv, "fake.echo"), "offline not delegated"); +} + +static void test_stale_heartbeat(kvlangKv_t *kv) { + register_backend(kv, "stale", "stale.op", "ready"); + char *hk = kvlangKeytreeSysRwirBackendHeartbeat("stale"); + set_char(kv, hk, "1"); + free(hk); + expect(!kvlangDispatchIsDelegatedOp(kv, "stale.op"), "stale heartbeat"); +} + +static void test_delegate_ok(kvlangKv_t *kv) { + register_backend(kv, "okb", "ok.echo", "ready"); + char *cmd = kvlangKeytreeSysRwirBackendCmd("okb"); + exec_arg_t a = { kv, cmd, "/data/out", NULL, 1, 0, 0 }; + pthread_t th; + pthread_create(&th, NULL, executor, &a); + kvlangRwirInst_t inst = make_inst("ok.echo", "/data/out"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + pthread_join(th, NULL); + expect(rc == KVLANG_DELEGATE_OK, "delegate ok"); + char *out = get_char(kv, "/data/out"); + expect(out && strcmp(out, "ok") == 0, "output written"); + free(out); + char **names = NULL; int n = 0; + kvlangKvList(kv, "/sys/task/", false, false, &names, &n); + expect(n == 0, "delegate must expire task status off the directory"); + for (int i = 0; i < n; i++) free(names[i]); + free(names); + kvlangRwirInstFree(&inst); + free(cmd); +} + +static void test_write_target(kvlangKv_t *kv) { + register_backend(kv, "wb", "w.echo", "ready"); + kvlangRwirInst_t inst = make_inst("w.echo", "/lib/evil"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + expect(rc == KVLANG_DELEGATE_ERR, "write-target fails"); + char *pc = NULL, *st = NULL; + kvlangVthreadGet(kv, "1", &pc, &st); + expect(st && strcmp(st, "error") == 0, "error status"); + free(pc); free(st); + kvlangRwirInstFree(&inst); +} + +static void test_done_without_output(kvlangKv_t *kv) { + register_backend(kv, "nb", "n.echo", "ready"); + char *cmd = kvlangKeytreeSysRwirBackendCmd("nb"); + exec_arg_t a = { kv, cmd, "/data/nout", NULL, 0, 0, 0 }; + pthread_t th; + pthread_create(&th, NULL, executor, &a); + kvlangRwirInst_t inst = make_inst("n.echo", "/data/nout"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + pthread_join(th, NULL); + expect(rc == KVLANG_DELEGATE_ERR, "empty output is failure"); + kvlangRwirInstFree(&inst); + free(cmd); +} + +static void test_timeout(kvlangKv_t *kv) { + register_backend(kv, "tb", "t.echo", "ready"); + kvlangDispatchSetDefaultTimeoutNs(40000000LL); + kvlangRwirInst_t inst = make_inst("t.echo", "/data/tout"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + kvlangDispatchSetDefaultTimeoutNs(30LL * 1000000000LL); + expect(rc == KVLANG_DELEGATE_ERR, "timeout"); + kvlangRwirInstFree(&inst); +} + +static char *error_msg(kvlangKv_t *kv, const char *vtid) { + kvlangStrbuf_t k; kvlangStrbufInit(&k); + kvlangKeytreeVthreadStatusMsg(vtid, "error", &k); + char *s = get_char(kv, k.p); + kvlangStrbufFree(&k); + return s; +} + +static void test_one_definition(kvlangKv_t *kv) { + register_backend(kv, "db", "dup.op", "ready"); + uint8_t raw[8] = { 0, 0, 0, 0, 0, 0, 0, 0 }; + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewTlv(&v, KVSPACE_KIND_RWFUNC, raw, 4, 1); + char *sk = kvlangKeytreeLibSig("dup.op"); + kvlangKvPair_t p = { sk, v }; + char err[256]; + kvlangKvSet(kv, &p, 1, err, sizeof err); + kvlangXvalueFree(&v); + kvlangRwirInst_t inst = make_inst("dup.op", "/data/dup"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + expect(rc == KVLANG_DELEGATE_ERR, "rwfunc + backend conflict"); + char *msg = error_msg(kv, "1"); + expect(msg && strstr(msg, "one opcode, one definition"), "conflict names both definitions"); + free(msg); + kvlangRwirInstFree(&inst); + free(sk); +} + +static void *executor_out_only(void *p) { + exec_arg_t *a = p; + char *json = take_char(a->kv, a->cmd, 2000000000ULL); + if (json) { + set_char(a->kv, a->out_key, "ok"); + free(json); + a->ran = 1; + } + return NULL; +} + +static void test_output_without_done(kvlangKv_t *kv) { + register_backend(kv, "od", "od.echo", "ready"); + char *cmd = kvlangKeytreeSysRwirBackendCmd("od"); + exec_arg_t a = { kv, cmd, "/data/odout", NULL, 1, 0, 0 }; + pthread_t th; + pthread_create(&th, NULL, executor_out_only, &a); + kvlangDispatchSetDefaultTimeoutNs(80000000LL); + kvlangRwirInst_t inst = make_inst("od.echo", "/data/odout"); + kvlangVthreadSet(kv, "1", "/vthread/1/[1,0]", "running"); + int rc = kvlangDispatchDelegate(kv, "1", "/vthread/1/[1,0]", &inst); + pthread_join(th, NULL); + kvlangDispatchSetDefaultTimeoutNs(30LL * 1000000000LL); + expect(rc == KVLANG_DELEGATE_ERR, "output without done is failure"); + kvlangRwirInstFree(&inst); + free(cmd); +} + +static void test_missing_heartbeat(kvlangKv_t *kv) { + char *dir = kvlangKeytreeSysRwirBackend("nohb"); + char path[256]; + snprintf(path, sizeof path, "%s/", dir); + char err[256]; + kvlangKvMkindex(kv, path, err, sizeof err); + snprintf(path, sizeof path, "%s/op/", dir); + kvlangKvMkindex(kv, path, err, sizeof err); + free(dir); + char *opk = kvlangKeytreeSysRwirBackendOp("nohb", "nohb.op"); + set_char(kv, opk, "1"); + free(opk); + char *sk = kvlangKeytreeSysRwirBackendStatus("nohb"); + set_char(kv, sk, "ready"); + free(sk); + expect(!kvlangDispatchIsDelegatedOp(kv, "nohb.op"), "missing heartbeat is off duty"); +} + +int main(void) { + test_check_write_key(); + kvlangKv_t *kv = open_kv(); + test_is_delegated(kv); + test_stale_heartbeat(kv); + test_delegate_ok(kv); + test_write_target(kv); + test_done_without_output(kv); + test_timeout(kv); + test_one_definition(kv); + test_output_without_done(kv); + test_missing_heartbeat(kv); + kvlangKvDisconnect(kv); + fprintf(stderr, "ok\n"); + return 0; +} diff --git a/runtime/tests/expire_test.c b/runtime/tests/expire_test.c new file mode 100644 index 00000000..c75f8f52 --- /dev/null +++ b/runtime/tests/expire_test.c @@ -0,0 +1,101 @@ +#include "runtime_internal.h" +#include + +static void fail(const char *msg) { + fprintf(stderr, "FAIL %s\n", msg); + exit(1); +} + +static void expect(bool ok, const char *msg) { + if (!ok) fail(msg); +} + +static kvlangKv_t *open_kv(void) { + char dsn[128]; + snprintf(dsn, sizeof dsn, "shm:///tmp/kvlang-expire-%d", (int)getpid()); + kvlangKv_t *kv = kvlangKvConnect(dsn); + if (!kv) fail("connect"); + char err[256]; + kvlangKvMkindex(kv, "/e/", err, sizeof err); + kvlangKvMkindex(kv, "/e2/", err, sizeof err); + kvlangKvMkindex(kv, "/e/kdir/", err, sizeof err); + return kv; +} + +static void set_char(kvlangKv_t *kv, const char *key, const char *s) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewCharUtf8(&v, s); + kvlangKvPair_t p = { (char *)key, v }; + char err[256]; + if (kvlangKvSet(kv, &p, 1, err, sizeof err) != 0) fail(err[0] ? err : "set"); + kvlangXvalueFree(&v); +} + +static char *get_char(kvlangKv_t *kv, const char *key) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangKvGetOne(kv, key, &v); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +static int list_has(kvlangKv_t *kv, const char *dir, const char *name) { + char **names = NULL; int n = 0; + kvlangKvList(kv, dir, false, false, &names, &n); + int hit = 0; + for (int i = 0; i < n; i++) { + if (strcmp(names[i], name) == 0) hit = 1; + free(names[i]); + } + free(names); + return hit; +} + +int main(void) { + kvlangKv_t *kv = open_kv(); + char err[256]; + + set_char(kv, "/e2/only", "v"); + expect(kvlangKvExpire(kv, "/e2/only", 0, err, sizeof err) != 0, "zero duration must fail"); + expect(kvlangKvExpire(kv, "relkey", 40000000ULL, err, sizeof err) != 0, "relative key must fail"); + expect(kvlangKvExpire(kv, "/e2/missing", 40000000ULL, err, sizeof err) != 0, "missing key must fail"); + expect(kvlangKvExpire(kv, "/e2/", 40000000ULL, err, sizeof err) != 0, "directory must fail"); + + set_char(kv, "/e/k", "v"); + set_char(kv, "/e/kdir/x", "1"); + expect(kvlangKvExpire(kv, "/e/k", 40000000ULL, err, sizeof err) == 0, "expire file"); + char *got = get_char(kv, "/e/k"); + expect(got && strcmp(got, "v") == 0, "Get must still see the value before TTL"); + free(got); + expect(!list_has(kv, "/e/", "k"), "List must drop the file name at Expire"); + expect(list_has(kv, "/e/", "kdir"), "Expire of a file must not drop the sibling dir"); + got = get_char(kv, "/e/kdir/x"); + expect(got && strcmp(got, "1") == 0, "sibling dir contents must survive Expire of the file"); + free(got); + + set_char(kv, "/e2/only2", "v"); + expect(kvlangKvExpire(kv, "/e2/only2", 40000000ULL, err, sizeof err) == 0, "expire only"); + expect(!list_has(kv, "/e2/", "only2"), "List must drop the lone file"); + + int gone = 0; + for (int i = 0; i < 40; i++) { + got = get_char(kv, "/e/k"); + if (!got) { gone = 1; break; } + free(got); + usleep(10000); + } + expect(gone, "Get must be None after TTL"); + + set_char(kv, "/e2/ghost", "g"); + expect(kvlangKvExpire(kv, "/e2/ghost", 40000000ULL, err, sizeof err) == 0, "expire ghost"); + usleep(80000); + expect(!list_has(kv, "/e2/", "ghost"), "List after TTL must not resurrect"); + expect(!list_has(kv, "/e2/", "ghost"), "second List after TTL must stay empty"); + got = get_char(kv, "/e2/ghost"); + expect(!got, "Get after List-first reap is None"); + free(got); + + kvlangKvDisconnect(kv); + fprintf(stderr, "ok\n"); + return 0; +} diff --git a/runtime/tests/incr_test.c b/runtime/tests/incr_test.c new file mode 100644 index 00000000..535c80b3 --- /dev/null +++ b/runtime/tests/incr_test.c @@ -0,0 +1,79 @@ +#include "runtime_internal.h" +#include +#include +#include + +#define NTHREAD 8 +#define NPER 100 + +static void fail(const char *msg) { + fprintf(stderr, "FAIL %s\n", msg); + exit(1); +} + +static kvlangKv_t *open_kv(void) { + char dsn[128]; + snprintf(dsn, sizeof dsn, "shm:///tmp/kvlang-incr-%d-%ld", (int)getpid(), (long)time(NULL)); + kvlangKv_t *kv = kvlangKvConnect(dsn); + if (!kv) fail("connect"); + return kv; +} + +typedef struct { kvlangKv_t *kv; const char *key; int64_t got[NPER]; } incr_arg_t; + +static void *worker(void *p) { + incr_arg_t *a = p; + for (int i = 0; i < NPER; i++) a->got[i] = kvlangVthreadNextSeq(a->kv, a->key); + return NULL; +} + +int main(void) { + kvlangKv_t *kv = open_kv(); + + int64_t first = 0; + char err[256]; + if (kvlangKvIncr(kv, "/seq-start", &first, err, sizeof err) != 0) fail("first incr"); + if (first != 1) fail("missing key must start at 1"); + for (int i = 2; i <= 12; i++) { + int64_t n = 0; + if (kvlangKvIncr(kv, "/seq-start", &n, err, sizeof err) != 0) fail("walk incr"); + if (n != i) fail("char counter must parse multi-digit values"); + } + + const char *key = "/seq-race"; + incr_arg_t args[NTHREAD]; + pthread_t th[NTHREAD]; + for (int i = 0; i < NTHREAD; i++) { + args[i].kv = kv; + args[i].key = key; + pthread_create(&th[i], NULL, worker, &args[i]); + } + for (int i = 0; i < NTHREAD; i++) pthread_join(th[i], NULL); + + int seen[NTHREAD * NPER + 1]; + memset(seen, 0, sizeof seen); + int max = 0; + for (int t = 0; t < NTHREAD; t++) { + for (int i = 0; i < NPER; i++) { + int64_t v = args[t].got[i]; + if (v < 1 || v > NTHREAD * NPER) fail("seq out of range"); + if (seen[v]) fail("duplicate seq"); + seen[v] = 1; + if ((int)v > max) max = (int)v; + } + } + if (max != NTHREAD * NPER) fail("missing seq values"); + + kvlangXvalue_t g; kvlangXvalueZero(&g); + kvlangKvGetOne(kv, key, &g); + char *s = kvlangXvalueNone(&g) ? NULL : kvlangXvalueValueString(&g); + kvlangXvalueFree(&g); + char want[32]; + snprintf(want, sizeof want, "%d", NTHREAD * NPER); + if (!s || strcmp(s, want) != 0) fail("persisted counter != last seq"); + free(s); + + kvlangKvDisconnect(kv); + fprintf(stderr, "ok\n"); + return 0; +} diff --git a/runtime/tests/notify_contract_test.c b/runtime/tests/notify_contract_test.c new file mode 100644 index 00000000..24b16de3 --- /dev/null +++ b/runtime/tests/notify_contract_test.c @@ -0,0 +1,76 @@ +#include "runtime_internal.h" +#include + +static void fail(const char *msg) { + fprintf(stderr, "FAIL %s\n", msg); + exit(1); +} + +static void expect(bool ok, const char *msg) { + if (!ok) fail(msg); +} + +static kvlangKv_t *open_kv(void) { + char dsn[128]; + snprintf(dsn, sizeof dsn, "shm:///tmp/kvlang-notify-%d", (int)getpid()); + kvlangKv_t *kv = kvlangKvConnect(dsn); + if (!kv) fail("connect"); + return kv; +} + +static void notify_char(kvlangKv_t *kv, const char *key, const char *s) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewCharUtf8(&v, s); + char err[256]; + if (kvlangKvNotify(kv, key, &v, err, sizeof err) != 0) fail(err[0] ? err : "notify"); + kvlangXvalueFree(&v); +} + +static char *take_char(kvlangKv_t *kv, const char *key, uint64_t ns) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + if (kvlangKvTake(kv, key, ns, &v) != 0) fail("take"); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +int main(void) { + kvlangKv_t *kv = open_kv(); + const uint64_t to = 200000000ULL; + + notify_char(kv, "/q/early", "1"); + char *got = take_char(kv, "/q/early", to); + expect(got && strcmp(got, "1") == 0, "a Notify posted before Take was lost"); + free(got); + + notify_char(kv, "/q/queued", "first"); + notify_char(kv, "/q/queued", "second"); + char *a = take_char(kv, "/q/queued", to); + char *b = take_char(kv, "/q/queued", to); + expect(a && b, "only one of 2 queued tasks was delivered"); + int saw_first = (strcmp(a, "first") == 0) || (strcmp(b, "first") == 0); + int saw_second = (strcmp(a, "second") == 0) || (strcmp(b, "second") == 0); + expect(saw_first && saw_second, "both tasks must survive"); + free(a); free(b); + + notify_char(kv, "/q/consumed", "x"); + got = take_char(kv, "/q/consumed", to); + expect(got && strcmp(got, "x") == 0, "the only Notify was not delivered"); + free(got); + got = take_char(kv, "/q/consumed", to); + expect(!got, "Take redelivered a consumed task"); + free(got); + + notify_char(kv, "/q/outside", "hidden"); + kvlangXvalue_t tree; kvlangXvalueZero(&tree); + kvlangKvGetOne(kv, "/q/outside", &tree); + expect(kvlangXvalueNone(&tree), "Notify must not write the user-visible key"); + kvlangXvalueFree(&tree); + got = take_char(kv, "/q/outside", to); + expect(got && strcmp(got, "hidden") == 0, "side queue lost the Notify"); + free(got); + + kvlangKvDisconnect(kv); + fprintf(stderr, "ok\n"); + return 0; +} diff --git a/runtime/tests/watchany_test.c b/runtime/tests/watchany_test.c new file mode 100644 index 00000000..c3891f56 --- /dev/null +++ b/runtime/tests/watchany_test.c @@ -0,0 +1,72 @@ +#include "runtime_internal.h" +#include + +static void fail(const char *msg) { + fprintf(stderr, "FAIL %s\n", msg); + exit(1); +} + +static void expect(bool ok, const char *msg) { + if (!ok) fail(msg); +} + +static kvlangKv_t *open_kv(void) { + char dsn[128]; + snprintf(dsn, sizeof dsn, "shm:///tmp/kvlang-watchany-%d", (int)getpid()); + kvlangKv_t *kv = kvlangKvConnect(dsn); + if (!kv) fail("connect"); + return kv; +} + +static void notify_char(kvlangKv_t *kv, const char *key, const char *s) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + kvlangXvalueNewCharUtf8(&v, s); + char err[256]; + if (kvlangKvNotify(kv, key, &v, err, sizeof err) != 0) fail(err[0] ? err : "notify"); + kvlangXvalueFree(&v); +} + +static char *take_char(kvlangKv_t *kv, const char *key, uint64_t ns) { + kvlangXvalue_t v; kvlangXvalueZero(&v); + if (kvlangKvTake(kv, key, ns, &v) != 0) fail("take"); + char *s = kvlangXvalueNone(&v) ? NULL : kvlangXvalueValueString(&v); + kvlangXvalueFree(&v); + return s; +} + +int main(void) { + kvlangKv_t *kv = open_kv(); + const uint64_t to = 200000000ULL; + const char *a = "/wa/a"; + const char *b = "/wa/b"; + const char *keys[2] = { a, b }; + + notify_char(kv, b, "B"); + char *got_key = NULL; + kvlangXvalue_t got; kvlangXvalueZero(&got); + if (kvlangKvWatchAny(kv, keys, 2, to, &got_key, &got) != 0) fail("watchany"); + char *got_s = kvlangXvalueNone(&got) ? NULL : kvlangXvalueValueString(&got); + expect(got_key && strcmp(got_key, b) == 0 && got_s && strcmp(got_s, "B") == 0, + "WatchAny must return the notified key"); + free(got_key); free(got_s); kvlangXvalueFree(&got); + + notify_char(kv, a, "A"); + notify_char(kv, b, "B2"); + got_key = NULL; kvlangXvalueZero(&got); + if (kvlangKvWatchAny(kv, keys, 2, to, &got_key, &got) != 0) fail("watchany2"); + expect(got_key && !kvlangXvalueNone(&got), "WatchAny lost a queued Notify"); + const char *other = (got_key && strcmp(got_key, a) == 0) ? b : a; + if (got_key && strcmp(got_key, a) != 0 && strcmp(got_key, b) != 0) fail("unexpected key"); + char *left = take_char(kv, other, to); + expect(left != NULL, "WatchAny must consume only the delivered key"); + free(left); free(got_key); kvlangXvalueFree(&got); + + got_key = NULL; kvlangXvalueZero(&got); + if (kvlangKvWatchAny(kv, keys, 2, to, &got_key, &got) != 0) fail("watchany empty"); + expect(!got_key && kvlangXvalueNone(&got), "WatchAny must time out when nothing is queued"); + kvlangXvalueFree(&got); + + kvlangKvDisconnect(kv); + fprintf(stderr, "ok\n"); + return 0; +}