feat: arc stage 3 T7 — transparent DB actor (WO_T_DB hole closed)
- worker DB builtins marshal to shard 0: requester-side slot encode (VM heaps never read cross-shard), owner executes serialized in adopt, reply unparks via new WO_PARK_INBOX park + envelope 3/4 - engine gains thread-agnostic slot entry points (insert_slots, update_field_slot, val_encode/clone, wo_db_exec_req); traps and messages byte-identical to the local path - main.c: engine + replay boot BEFORE shards spawn; workers assert rt.db/rt.wal NULL; busy shard adopts inbox once per slice - latent stage-1 bug fixed: shared io_uring params static raced by lazy worker init lost park wakes (~1/20 hangs); params per-vm, short submit now fails loud - new sample docs/examples/db-actor + just db-actor gate 8/0 (multi x3, uring/epoll forced, single byte-exact, WAL replay pair); ASan+TSan 6/6; full battery green Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
954aacc9a5
commit
1ed4922fdd
15 changed files with 911 additions and 47 deletions
|
|
@ -1,5 +1,6 @@
|
|||
#include "db.h"
|
||||
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "cont.h"
|
||||
|
|
@ -161,3 +162,176 @@ int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
|||
return WO_T_DB;
|
||||
}
|
||||
}
|
||||
|
||||
/* ---- arc stage 3: the owner-shard executor ------------------------------
|
||||
* Mirrors the switch above case for case, with slot inputs and plain
|
||||
* outputs — every trap code and message a worker sees is byte-identical to
|
||||
* what the same statement would produce on the primary. */
|
||||
void wo_db_exec_req(wo_vm *vm, wo_db_req *q) {
|
||||
wo_db *db = (wo_db *)vm->rt.db;
|
||||
wo_wal *w = (wo_wal *)vm->rt.wal;
|
||||
const char *m = "db failed";
|
||||
q->status = 0;
|
||||
q->msg = "";
|
||||
if (!db) {
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "database engine not initialized";
|
||||
goto out;
|
||||
}
|
||||
switch (q->op) {
|
||||
case WO_B_DB_INSERT: {
|
||||
int ek = 0;
|
||||
uint64_t id = wo_row_insert_slots(db, q->cid, q->slots, &m, &ek);
|
||||
if (!id) {
|
||||
q->status = ek == DB_ERR_UNIQUE ? WO_T_UNIQUE
|
||||
: ek == DB_ERR_OOM ? WO_T_OOM
|
||||
: WO_T_DB;
|
||||
q->msg = m;
|
||||
break;
|
||||
}
|
||||
if (w) {
|
||||
if (wo_wal_append_insert(w, db, q->cid, id) != 0 || wo_wal_commit(w) != 0) {
|
||||
wo_row_remove(db, q->cid, id);
|
||||
q->status = WO_T_IO;
|
||||
q->msg = "wal commit failed";
|
||||
break;
|
||||
}
|
||||
}
|
||||
q->result = id;
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_UPDATE_FIELD: {
|
||||
int ek = 0;
|
||||
if (wo_row_update_field_slot(db, q->cid, q->id, q->field, q->slots[0], &m, &ek) != 0) {
|
||||
q->status = ek == DB_ERR_UNIQUE ? WO_T_UNIQUE : ek == DB_ERR_OOM ? WO_T_OOM : WO_T_DB;
|
||||
q->msg = m;
|
||||
break;
|
||||
}
|
||||
if (w) {
|
||||
if (wo_wal_append_update(w, db, q->cid, q->id) != 0 || wo_wal_commit(w) != 0) {
|
||||
q->status = WO_T_IO;
|
||||
q->msg = "wal commit failed";
|
||||
break;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_DELETE: {
|
||||
if (wo_row_has_referrers(db, q->cid, q->id)) {
|
||||
q->status = WO_T_FK;
|
||||
q->msg = "row is still referenced (restrict)";
|
||||
break;
|
||||
}
|
||||
if (wo_row_remove(db, q->cid, q->id) != 0) {
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "no such row";
|
||||
break;
|
||||
}
|
||||
if (w) {
|
||||
if (wo_wal_append_remove(w, q->cid, q->id) != 0 || wo_wal_commit(w) != 0) {
|
||||
q->status = WO_T_IO;
|
||||
q->msg = "wal commit failed";
|
||||
break;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_SCAN:
|
||||
case WO_B_DB_PROBE: {
|
||||
if (q->cid >= db->class_cnt) {
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "no such class";
|
||||
break;
|
||||
}
|
||||
db_table *t = &db->tables[q->cid];
|
||||
uint64_t *out = NULL;
|
||||
uint32_t n = 0, cap = 0;
|
||||
if (t->row_size && (q->op == WO_B_DB_SCAN || q->index < t->index_cnt)) {
|
||||
uint32_t col = 0;
|
||||
uint8_t kind = 0;
|
||||
if (q->op == WO_B_DB_PROBE) {
|
||||
col = t->indexes[q->index].cols[0];
|
||||
kind = db->classes[q->cid].kinds[col];
|
||||
}
|
||||
uint32_t total = t->slab_cnt * DB_SLAB_ROWS;
|
||||
for (uint32_t g = 0; g < total; g++) {
|
||||
if (!(t->bitmap[g >> 6] & (1ull << (g & 63)))) continue;
|
||||
db_row *row =
|
||||
(db_row *)(t->slabs[g / DB_SLAB_ROWS] + (size_t)(g % DB_SLAB_ROWS) * t->row_size);
|
||||
if (q->op == WO_B_DB_PROBE) {
|
||||
int eq;
|
||||
if (kind == WO_K_TEXT || kind == WO_K_BYTES) {
|
||||
/* both sides engine-encoded: the key was encoded on
|
||||
* the requester's thread, the slot lives here */
|
||||
const db_text *want = (const db_text *)(uintptr_t)q->slots[0];
|
||||
const db_text *have = (const db_text *)(uintptr_t)row->slots[col];
|
||||
eq = (!want && !have) ||
|
||||
(want && have && want->len == have->len &&
|
||||
memcmp(want->bytes, have->bytes, have->len) == 0);
|
||||
} else
|
||||
eq = row->slots[col] == q->slots[0];
|
||||
if (!eq) continue;
|
||||
}
|
||||
if (n == cap) {
|
||||
uint32_t ncap = cap ? cap * 2 : 16;
|
||||
uint64_t *no = realloc(out, (size_t)ncap * 8u);
|
||||
if (!no) {
|
||||
free(out);
|
||||
out = NULL;
|
||||
q->status = WO_T_OOM;
|
||||
q->msg = "out of memory";
|
||||
break;
|
||||
}
|
||||
out = no;
|
||||
cap = ncap;
|
||||
}
|
||||
out[n++] = row->id;
|
||||
}
|
||||
}
|
||||
if (!q->status) {
|
||||
q->ids = out;
|
||||
q->id_cnt = n;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_GET_FIELD: {
|
||||
if (q->cid >= db->class_cnt || q->field >= db->classes[q->cid].field_cnt) {
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "no such field";
|
||||
break;
|
||||
}
|
||||
db_row *row = wo_row_ptr(db, q->cid, q->id);
|
||||
if (!row) {
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "no such row";
|
||||
break;
|
||||
}
|
||||
int ok = 1;
|
||||
q->val_kind = db->classes[q->cid].kinds[q->field];
|
||||
q->val = wo_db_val_clone(db->classes, q->val_kind, row->slots[q->field], &ok);
|
||||
if (!ok) {
|
||||
q->status = WO_T_OOM;
|
||||
q->msg = "out of memory";
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
q->status = WO_T_DB;
|
||||
q->msg = "unknown db builtin";
|
||||
break;
|
||||
}
|
||||
out:
|
||||
/* slot VALUES were consumed by the ops above (insert/update install or
|
||||
* free them); the PROBE key is ours to free, the array always is. (The
|
||||
* requester only encodes a key for an index its identical class table
|
||||
* declares, so a keyed request always finds its kind here.) */
|
||||
if (db && q->op == WO_B_DB_PROBE && q->slots && q->cid < db->class_cnt) {
|
||||
db_table *t = &db->tables[q->cid];
|
||||
if (t->row_size && q->index < t->index_cnt)
|
||||
wo_db_val_free(db, db->classes[q->cid].kinds[t->indexes[q->index].cols[0]],
|
||||
q->slots[0]);
|
||||
}
|
||||
free(q->slots);
|
||||
q->slots = NULL;
|
||||
q->slot_cnt = 0;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -21,4 +21,36 @@
|
|||
|
||||
int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg);
|
||||
|
||||
/* ---- arc stage 3: one marshaled DB statement (the transparent DB actor).
|
||||
* A worker shard fills the request on ITS thread — args pre-encoded into
|
||||
* engine slots, since VM heaps are never read cross-shard — and ships it
|
||||
* to shard 0 in an envelope; the owner executes it via wo_db_exec_req and
|
||||
* ships it back. Ownership: `slots` VALUES pass to the owner (consumed by
|
||||
* the op), the array and every reply buffer pass back to the requester.
|
||||
* The envelope handoff (mutex + eventfd) orders `done` on both sides. */
|
||||
typedef struct wo_db_req {
|
||||
/* request */
|
||||
uint32_t op; /* WO_B_DB_INSERT..WO_B_DB_PROBE */
|
||||
uint32_t cid, field, index;
|
||||
uint64_t id;
|
||||
uint64_t *slots; /* INSERT: field_cnt; UPDATE: 1; PROBE: 1 (the key) */
|
||||
uint32_t slot_cnt;
|
||||
/* routing */
|
||||
uint32_t from_shard;
|
||||
void *fiber; /* the parked wo_fiber*, opaque to the engine */
|
||||
int done;
|
||||
/* reply */
|
||||
int status; /* 0 ok, else the WO_T_* the local path would trap */
|
||||
const char *msg; /* static literal, safe cross-thread */
|
||||
uint64_t result; /* INSERT: the new id */
|
||||
uint64_t *ids; /* SCAN/PROBE: malloc'd id list */
|
||||
uint32_t id_cnt;
|
||||
uint8_t val_kind; /* GET_FIELD: a cloned engine value */
|
||||
uint64_t val;
|
||||
} wo_db_req;
|
||||
|
||||
/* Execute one marshaled statement on the OWNER shard (vm = shard 0's; its
|
||||
* rt.db/rt.wal are the engine). Fills the reply fields; never traps. */
|
||||
void wo_db_exec_req(wo_vm *vm, wo_db_req *q);
|
||||
|
||||
#endif /* WO_DB_H */
|
||||
|
|
|
|||
|
|
@ -591,6 +591,56 @@ uint64_t wo_row_insert(wo_db *db, uint32_t class_id, const uint64_t *vals,
|
|||
return r->id;
|
||||
}
|
||||
|
||||
uint64_t wo_row_insert_slots(wo_db *db, uint32_t class_id, const uint64_t *slots,
|
||||
const char **msg, int *err_kind) {
|
||||
if (err_kind) *err_kind = DB_ERR_MISC;
|
||||
db_table *t = table_of(db, class_id);
|
||||
const wo_classdesc *c = class_id < db->class_cnt ? &db->classes[class_id] : NULL;
|
||||
if (!t || !c) {
|
||||
/* slot kinds unknowable without the class: the values leak rather
|
||||
than die by the wrong kind (defensive; the requester validated) */
|
||||
*msg = "no such class";
|
||||
return 0;
|
||||
}
|
||||
uint32_t g = slot_alloc(t);
|
||||
if (g == UINT32_MAX) {
|
||||
for (uint32_t j = 0; j < c->field_cnt; j++) db_val_free(c->kinds[j], slots[j]);
|
||||
if (err_kind) *err_kind = DB_ERR_OOM;
|
||||
*msg = "out of memory growing a table";
|
||||
return 0;
|
||||
}
|
||||
db_row *r = slot_row(t, g);
|
||||
r->class_id = class_id;
|
||||
r->flags = 0;
|
||||
memcpy(r->slots, slots, (size_t)c->field_cnt * 8u);
|
||||
r->id = t->next_id;
|
||||
t->next_id += db->nshards;
|
||||
if (hput(t, r->id, (uint64_t)g + 1) != 0) {
|
||||
for (uint32_t j = 0; j < c->field_cnt; j++) db_val_free(c->kinds[j], r->slots[j]);
|
||||
t->next_id -= db->nshards;
|
||||
if (t->free_cnt < t->free_cap) t->free_slots[t->free_cnt++] = g;
|
||||
if (err_kind) *err_kind = DB_ERR_OOM;
|
||||
*msg = "out of memory indexing a row";
|
||||
return 0;
|
||||
}
|
||||
t->bitmap[g >> 6] |= 1ull << (g & 63);
|
||||
t->count++;
|
||||
int irc = idx_add_row(db, t, r);
|
||||
if (irc != 0) {
|
||||
t->bitmap[g >> 6] &= ~(1ull << (g & 63));
|
||||
hdel(t, r->id);
|
||||
t->count--;
|
||||
t->next_id -= db->nshards; /* the id was never observable: reclaim it */
|
||||
for (uint32_t j = 0; j < c->field_cnt; j++) db_val_free(c->kinds[j], r->slots[j]);
|
||||
if (t->free_cnt < t->free_cap) t->free_slots[t->free_cnt++] = g;
|
||||
if (err_kind) *err_kind = irc;
|
||||
*msg = irc == DB_ERR_UNIQUE ? "unique index violation" : "out of memory indexing a row";
|
||||
return 0;
|
||||
}
|
||||
if (err_kind) *err_kind = DB_ERR_NONE;
|
||||
return r->id;
|
||||
}
|
||||
|
||||
db_row *wo_row_ptr(wo_db *db, uint32_t class_id, uint64_t id) {
|
||||
if (class_id >= db->class_cnt) return NULL;
|
||||
db_table *t = &db->tables[class_id];
|
||||
|
|
@ -645,12 +695,101 @@ void wo_db_val_free(wo_db *db, uint8_t kind, uint64_t v) {
|
|||
db_val_free(kind, v);
|
||||
}
|
||||
|
||||
uint64_t wo_db_val_encode(const wo_classdesc *classes, uint8_t kind, uint64_t vm_val,
|
||||
int *ok, const char **msg) {
|
||||
return db_val_encode(classes, kind, vm_val, ok, msg);
|
||||
}
|
||||
|
||||
/* Deep engine-to-engine copy; shapes mirror db_val_free's recursion. */
|
||||
uint64_t wo_db_val_clone(const wo_classdesc *classes, uint8_t kind, uint64_t v, int *ok) {
|
||||
*ok = 1;
|
||||
if (!v) return 0;
|
||||
switch (kind) {
|
||||
case WO_K_SCALAR:
|
||||
case WO_K_FLOAT: return v;
|
||||
case WO_K_TEXT:
|
||||
case WO_K_BYTES: {
|
||||
const db_text *s = (const db_text *)(uintptr_t)v;
|
||||
db_text *t = malloc(sizeof(db_text) + s->len);
|
||||
if (!t) goto oom;
|
||||
t->len = s->len;
|
||||
memcpy(t->bytes, s->bytes, s->len);
|
||||
return (uint64_t)(uintptr_t)t;
|
||||
}
|
||||
case WO_K_OWNED: {
|
||||
const db_rec *s = (const db_rec *)(uintptr_t)v;
|
||||
const wo_classdesc *c = &classes[s->class_id];
|
||||
db_rec *r = malloc(sizeof(db_rec) + (size_t)c->field_cnt * 8u);
|
||||
if (!r) goto oom;
|
||||
r->class_id = s->class_id;
|
||||
r->_pad = 0;
|
||||
for (uint32_t i = 0; i < c->field_cnt; i++) {
|
||||
r->slots[i] = wo_db_val_clone(classes, c->kinds[i], s->slots[i], ok);
|
||||
if (!*ok) {
|
||||
for (uint32_t j = 0; j < i; j++) db_val_free(c->kinds[j], r->slots[j]);
|
||||
free(r);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
return (uint64_t)(uintptr_t)r;
|
||||
}
|
||||
case WO_K_MULTI: {
|
||||
const db_multi *s = (const db_multi *)(uintptr_t)v;
|
||||
db_multi *d = malloc(sizeof(db_multi) + (size_t)s->len * 8u);
|
||||
if (!d) goto oom;
|
||||
d->elem_kind = s->elem_kind;
|
||||
d->len = s->len;
|
||||
for (uint32_t i = 0; i < s->len; i++) {
|
||||
d->items[i] = wo_db_val_clone(classes, s->elem_kind, s->items[i], ok);
|
||||
if (!*ok) {
|
||||
for (uint32_t j = 0; j < i; j++) db_val_free(d->elem_kind, d->items[j]);
|
||||
free(d);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
return (uint64_t)(uintptr_t)d;
|
||||
}
|
||||
case WO_K_MAP: {
|
||||
const db_map *s = (const db_map *)(uintptr_t)v;
|
||||
db_map *d = malloc(sizeof(db_map) + (size_t)s->len * 16u);
|
||||
if (!d) goto oom;
|
||||
d->key_kind = s->key_kind;
|
||||
d->val_kind = s->val_kind;
|
||||
d->len = s->len;
|
||||
for (uint32_t i = 0; i < s->len; i++) {
|
||||
d->kv[2 * i] = wo_db_val_clone(classes, s->key_kind, s->kv[2 * i], ok);
|
||||
uint64_t dv = 0;
|
||||
if (*ok) dv = wo_db_val_clone(classes, s->val_kind, s->kv[2 * i + 1], ok);
|
||||
d->kv[2 * i + 1] = dv;
|
||||
if (!*ok) {
|
||||
for (uint32_t j = 0; j <= i; j++) {
|
||||
db_val_free(d->key_kind, d->kv[2 * j]);
|
||||
db_val_free(d->val_kind, d->kv[2 * j + 1]);
|
||||
}
|
||||
free(d);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
return (uint64_t)(uintptr_t)d;
|
||||
}
|
||||
default: return v; /* GCREF never stored; nothing to clone */
|
||||
}
|
||||
oom:
|
||||
*ok = 0;
|
||||
return 0;
|
||||
}
|
||||
|
||||
uint64_t wo_val_decode_vm(wo_db *db, wo_rt *rt, uint8_t kind, uint64_t engine_val,
|
||||
int *ok, const char **msg) {
|
||||
(void)db;
|
||||
return db_val_decode(rt, kind, engine_val, ok, msg);
|
||||
}
|
||||
|
||||
static int row_apply_field_slot(wo_db *db, db_table *t, const wo_classdesc *c,
|
||||
db_row *r, uint32_t class_id, uint64_t id,
|
||||
uint32_t field, uint64_t nv, const char **msg,
|
||||
int *err_kind);
|
||||
|
||||
int wo_row_update_field(wo_db *db, uint32_t class_id, uint64_t id, uint32_t field,
|
||||
uint64_t vm_val, const char **msg, int *err_kind) {
|
||||
if (err_kind) *err_kind = DB_ERR_MISC;
|
||||
|
|
@ -671,6 +810,16 @@ int wo_row_update_field(wo_db *db, uint32_t class_id, uint64_t id, uint32_t fiel
|
|||
if (err_kind) *err_kind = DB_ERR_BADKIND;
|
||||
return -1;
|
||||
}
|
||||
return row_apply_field_slot(db, t, c, r, class_id, id, field, nv, msg, err_kind);
|
||||
}
|
||||
|
||||
/* The post-encode half of an update: unique shadow-check, index fix-up,
|
||||
* slot swap. Consumes [nv] (installed on success, freed on failure) —
|
||||
* shared by the VM-value wrapper above and the RPC slot path. */
|
||||
static int row_apply_field_slot(wo_db *db, db_table *t, const wo_classdesc *c,
|
||||
db_row *r, uint32_t class_id, uint64_t id,
|
||||
uint32_t field, uint64_t nv, const char **msg,
|
||||
int *err_kind) {
|
||||
/* indexes containing this column: unique checks against the NEW value
|
||||
run first, against a shadow of the row, before anything mutates */
|
||||
uint64_t old = r->slots[field];
|
||||
|
|
@ -738,6 +887,31 @@ int wo_row_update_field(wo_db *db, uint32_t class_id, uint64_t id, uint32_t fiel
|
|||
return 0;
|
||||
}
|
||||
|
||||
int wo_row_update_field_slot(wo_db *db, uint32_t class_id, uint64_t id, uint32_t field,
|
||||
uint64_t slot, const char **msg, int *err_kind) {
|
||||
if (err_kind) *err_kind = DB_ERR_MISC;
|
||||
/* bounds first: the RPC requester validated cid/field to encode at all,
|
||||
so these are defensive; the slot's kind is unknowable on a class
|
||||
violation and the value leaks rather than dies by the wrong kind */
|
||||
if (class_id >= db->class_cnt) {
|
||||
*msg = "no such class";
|
||||
return -1;
|
||||
}
|
||||
const wo_classdesc *c = &db->classes[class_id];
|
||||
if (field >= c->field_cnt) {
|
||||
*msg = "no such field";
|
||||
return -1;
|
||||
}
|
||||
db_row *r = wo_row_ptr(db, class_id, id);
|
||||
if (!r) {
|
||||
db_val_free(c->kinds[field], slot);
|
||||
*msg = "no such row";
|
||||
return -1;
|
||||
}
|
||||
return row_apply_field_slot(db, &db->tables[class_id], c, r, class_id, id,
|
||||
field, slot, msg, err_kind);
|
||||
}
|
||||
|
||||
int wo_row_has_referrers(wo_db *db, uint32_t class_id, uint64_t id) {
|
||||
if (!id) return 0;
|
||||
for (uint32_t c = 0; c < db->class_cnt; c++) {
|
||||
|
|
|
|||
|
|
@ -193,4 +193,32 @@ uint64_t wo_val_decode_vm(wo_db *db, wo_rt *rt, uint8_t kind, uint64_t engine_va
|
|||
* (0 ok, -1). */
|
||||
int wo_row_raw_commit(wo_db *db, uint32_t class_id, db_row *r);
|
||||
|
||||
/* ---- arc stage 3: the slot-level surface the transparent DB RPC uses ----
|
||||
* VM heaps are never read cross-shard (a worker's GC writes header mark
|
||||
* bits concurrently), so the REQUESTER shard encodes its VM values into
|
||||
* engine-owned slots on its own thread and ships those; the OWNER shard
|
||||
* executes from slots. Everything here is thread-agnostic: it touches only
|
||||
* the wo_db it is handed and engine-owned mallocs. */
|
||||
|
||||
/* Encode one VM value into an engine slot on the caller's thread (the
|
||||
* in-gate, split out of wo_row_insert for the RPC path). */
|
||||
uint64_t wo_db_val_encode(const wo_classdesc *classes, uint8_t kind, uint64_t vm_val,
|
||||
int *ok, const char **msg);
|
||||
|
||||
/* Deep-copy one engine value — a GET_FIELD reply must outlive the row it
|
||||
* was read from (a later statement may free the row's slot). */
|
||||
uint64_t wo_db_val_clone(const wo_classdesc *classes, uint8_t kind, uint64_t v, int *ok);
|
||||
|
||||
/* Insert from PRE-ENCODED slots (field_cnt of them). Ownership of the slot
|
||||
* VALUES transfers: installed on success, freed on failure. New id, or 0
|
||||
* with *msg / *err_kind set exactly as wo_row_insert sets them. */
|
||||
uint64_t wo_row_insert_slots(wo_db *db, uint32_t class_id, const uint64_t *slots,
|
||||
const char **msg, int *err_kind);
|
||||
|
||||
/* Update one field from a PRE-ENCODED slot value (consumed either way:
|
||||
* installed on success, freed on failure). Same contract as
|
||||
* wo_row_update_field after its encode. */
|
||||
int wo_row_update_field_slot(wo_db *db, uint32_t class_id, uint64_t id, uint32_t field,
|
||||
uint64_t slot, const char **msg, int *err_kind);
|
||||
|
||||
#endif /* WO_TABLE_H */
|
||||
|
|
|
|||
62
docs/examples/db-actor/main.wo
Normal file
62
docs/examples/db-actor/main.wo
Normal file
|
|
@ -0,0 +1,62 @@
|
|||
use time
|
||||
|
||||
-- db-actor — arc stage 3's acceptance workload: the database is an
|
||||
-- actor on the owner shard (shard 0); a spawned actor placed on ANY
|
||||
-- shard reads and writes it through transparent RPC. Before stage 3 a
|
||||
-- worker-shard insert traps WO_T_DB ("database engine not
|
||||
-- initialized"); after, this program's output is shard-placement-
|
||||
-- independent: two writer lines and one exact main line.
|
||||
|
||||
@table(name: "notes", index: [tag])
|
||||
class Note {
|
||||
tag: Text
|
||||
val: Int
|
||||
}
|
||||
|
||||
class Job {
|
||||
n: Int
|
||||
}
|
||||
|
||||
-- Each writer inserts one row, then scans the whole table. Placement is
|
||||
-- round-robin, so with two writers at default shards one lands off the
|
||||
-- primary — the RPC path under test.
|
||||
class Writer {
|
||||
pad: Int
|
||||
fn receive(msg: Job) {
|
||||
insert Note { tag: "w", val: msg.n };
|
||||
let total = 0;
|
||||
for x in from n in Note select n {
|
||||
total = total + x.val;
|
||||
}
|
||||
print("writer ${msg.n} sees sum ${total}");
|
||||
}
|
||||
}
|
||||
|
||||
fn main() -> Int {
|
||||
let a: actor Job = spawn Writer { pad: 0 };
|
||||
let b: actor Job = spawn Writer { pad: 1 };
|
||||
send(a, Job { n: 1 });
|
||||
send(b, Job { n: 2 });
|
||||
-- no request/response surface yet (iteration 31): poll until both rows
|
||||
-- landed, then give the writers' own prints a beat before main returns
|
||||
-- (main-return reaps every other fiber, mid-print included)
|
||||
let tries = 0;
|
||||
let count = 0;
|
||||
while count < 2 and tries < 200 {
|
||||
time.sleep(10);
|
||||
count = 0;
|
||||
for x in from n in Note select n {
|
||||
count = count + 1;
|
||||
}
|
||||
tries = tries + 1;
|
||||
}
|
||||
time.sleep(1000);
|
||||
count = 0;
|
||||
let total = 0;
|
||||
for x in from n in Note select n {
|
||||
count = count + 1;
|
||||
total = total + x.val;
|
||||
}
|
||||
print("main sees ${count} rows, sum ${total}");
|
||||
return 0;
|
||||
}
|
||||
6
docs/examples/db-actor/wo.toml
Normal file
6
docs/examples/db-actor/wo.toml
Normal file
|
|
@ -0,0 +1,6 @@
|
|||
name = "db-actor"
|
||||
version = "0.1.0"
|
||||
description = "arc stage 3 proof: actors on worker shards read and write the database through the DB actor"
|
||||
|
||||
[runtime]
|
||||
wo = ">= 0.1"
|
||||
|
|
@ -203,12 +203,35 @@ transitively-traced check), corpus + TSan.
|
|||
0), `database/src/` untouched (the engine never learns), `runtime/src/vm.c`
|
||||
(request/reply parking).
|
||||
|
||||
- [ ] Non-owner DB builtins marshal statement + args to shard 0, park,
|
||||
- [x] Non-owner DB builtins marshal statement + args to shard 0, park,
|
||||
resume with materialized reply; `transaction { }` travels as one unit
|
||||
(18's staged batch stays owner-side).
|
||||
- [ ] `just employee` + `just web-app` at default cores, answers
|
||||
byte-identical to N=1 (iteration 8's criterion 4, the arc's headline
|
||||
proof). Commit.
|
||||
(18's staged batch stays owner-side — nothing to do until 18 unholds).
|
||||
DEVIATIONS, disclosed: (1) "database/src untouched" bent to
|
||||
"database/src gains thread-agnostic slot-level entry points"
|
||||
(wo_db_val_encode/clone, wo_row_insert_slots, wo_row_update_field_slot,
|
||||
wo_db_exec_req) — the owner thread must never read a requester's VM
|
||||
heap (concurrent mark-bit writes = TSan race), so the REQUESTER encodes
|
||||
args to engine slots and the owner executes from slots, replay-style;
|
||||
(2) the reply park is a new plane-less park (`WO_PARK_INBOX`), woken by
|
||||
the DB_RESP envelope (envelope kinds 3/4; `wo_io_unpark` exported);
|
||||
resume re-executes the builtin, which consumes the reply; (3) a busy
|
||||
shard adopts its inbox once per reduction slice, bounding a request's
|
||||
wait on a computing primary; (4) main.c boots the engine + replay
|
||||
BEFORE `wo_engine_start` (the replay-before-serve obligation — it also
|
||||
publishes the engine's class table to worker threads by the spawn);
|
||||
(5) EN ROUTE, a latent stage-1 bug fixed: io_uring ring params were ONE
|
||||
file static, rewritten by every shard's lazy init while other shards
|
||||
read offsets from it — submits landed at garbage offsets and parked
|
||||
fibers lost wakes (~1/20 hangs at default cores). Params now live
|
||||
per-vm (`io_params`), and a short `io_uring_enter` submit is a loud
|
||||
trap, never a success.
|
||||
- [x] Verified: `just db-actor` (NEW gate, 8/0 — worker-shard actors
|
||||
insert/scan/get through the DB actor; multi-shard set-asserted ×3 +
|
||||
both forced backends; single-shard byte-exact; WO_DATA pair proves a
|
||||
worker's write is ack-after-durable and replays). ASan 6/6 and TSan
|
||||
6/6 clean on the RPC path; 60/60 hang-free at default cores.
|
||||
`just employee` + `just web-app` at default cores green (byte-identical
|
||||
to N=1). Commit.
|
||||
|
||||
### Task 8 — the arc's closeout
|
||||
|
||||
|
|
|
|||
6
justfile
6
justfile
|
|
@ -54,6 +54,12 @@ web-app:
|
|||
fibers:
|
||||
./scripts/fibers-accept.sh
|
||||
|
||||
# db-actor: arc stage 3's gate (docs/examples/db-actor) — worker-shard
|
||||
# actors read/write the database through the transparent DB actor; WAL
|
||||
# replay pair included. `just db-actor` runs it.
|
||||
db-actor:
|
||||
./scripts/db-actor-accept.sh
|
||||
|
||||
# install-accept: extract the dist tarball to a temp prefix, PATH it, and prove
|
||||
# `woc version` + a from-scratch project build+run (self-located wovm) + the
|
||||
# wo-constraint refusal all work — the "tarball install actually works" gate.
|
||||
|
|
|
|||
|
|
@ -171,7 +171,15 @@ int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
|||
if (C == WO_B_JSON_ENCODE || C == WO_B_JSON_DECODE)
|
||||
return wo_builtin_json(vm, R, ins, msg);
|
||||
if (C >= WO_B_SYS_FIRST && C <= WO_B_PROC_RUN) return wo_builtin_sys(vm, R, ins, msg);
|
||||
if (C >= WO_B_DB_INSERT && C <= WO_B_DB_PROBE) return wo_builtin_db(vm, R, ins, msg);
|
||||
if (C >= WO_B_DB_INSERT && C <= WO_B_DB_PROBE) {
|
||||
/* arc stage 3: the database is an actor on shard 0. A worker shard
|
||||
* has no engine by design — its statement marshals, parks, resumes
|
||||
* with the materialized reply. The primary (and every single-shard
|
||||
* or test build) keeps the direct path bit for bit. */
|
||||
if (!vm->rt.db && !vm->is_primary && wo_eng.nshards > 1)
|
||||
return wo_db_rpc(vm, R, ins, msg);
|
||||
return wo_builtin_db(vm, R, ins, msg);
|
||||
}
|
||||
switch (C) {
|
||||
case WO_B_SPAWN: {
|
||||
/* arc: R[B] = the moved-in instance, R[B+1] = receive's method
|
||||
|
|
|
|||
|
|
@ -178,31 +178,17 @@ int main(int argc, char **argv) {
|
|||
return 2;
|
||||
}
|
||||
wo_tls_set(&VM);
|
||||
/* the arc's stage 2: all cores by default (the brave landing), one
|
||||
* pinned worker vm per extra core; WO_SHARDS caps or forces it */
|
||||
{
|
||||
long cores = sysconf(_SC_NPROCESSORS_ONLN);
|
||||
uint32_t nshards = cores > 0 ? (uint32_t)cores : 1;
|
||||
const char *se = getenv("WO_SHARDS");
|
||||
if (se && se[0]) {
|
||||
unsigned long v = strtoul(se, NULL, 10);
|
||||
if (v >= 1 && v <= WO_MAX_SHARDS) nshards = (uint32_t)v;
|
||||
}
|
||||
if (nshards > WO_MAX_SHARDS) nshards = WO_MAX_SHARDS;
|
||||
wo_eng.shards = SHARDS;
|
||||
if (wo_engine_start(&mod, heap_mb << 20, nshards) != 0) {
|
||||
fprintf(stderr, "wovm: cannot start %u shards\n", nshards);
|
||||
wo_engine_stop();
|
||||
wo_vm_destroy(&VM);
|
||||
wo_module_free(&mod);
|
||||
return 2;
|
||||
}
|
||||
}
|
||||
/* The database engine boots with the VM: every class IS a table.
|
||||
* Durability is opt-in — WO_DATA=<dir> opens <dir>/shard-0.wal,
|
||||
* replays it before the entry runs (boot-before-listeners doctrine),
|
||||
* and every insert commits before it acknowledges. Without WO_DATA
|
||||
* the engine runs RAM-only, which is what the corpus expects. */
|
||||
* the engine runs RAM-only, which is what the corpus expects.
|
||||
* Arc stage 3 obligation: this whole block runs BEFORE the worker
|
||||
* shards spawn — replay completes before anything can serve, and the
|
||||
* engine's immutable class-table pointer is published to the worker
|
||||
* threads by the spawn itself. The engine and WAL stay the PRIMARY's
|
||||
* alone (rt.db/rt.wal are never set on a worker); workers reach them
|
||||
* through the DB actor's message path. */
|
||||
if (wo_db_init(&DB, mod.classes, mod.class_cnt, 0, 1) != 0) {
|
||||
fprintf(stderr, "wovm: cannot initialize the database engine\n");
|
||||
wo_vm_destroy(&VM);
|
||||
|
|
@ -230,6 +216,28 @@ int main(int argc, char **argv) {
|
|||
}
|
||||
VM.rt.wal = &WAL;
|
||||
}
|
||||
/* the arc's stage 2: all cores by default (the brave landing), one
|
||||
* pinned worker vm per extra core; WO_SHARDS caps or forces it */
|
||||
{
|
||||
long cores = sysconf(_SC_NPROCESSORS_ONLN);
|
||||
uint32_t nshards = cores > 0 ? (uint32_t)cores : 1;
|
||||
const char *se = getenv("WO_SHARDS");
|
||||
if (se && se[0]) {
|
||||
unsigned long v = strtoul(se, NULL, 10);
|
||||
if (v >= 1 && v <= WO_MAX_SHARDS) nshards = (uint32_t)v;
|
||||
}
|
||||
if (nshards > WO_MAX_SHARDS) nshards = WO_MAX_SHARDS;
|
||||
wo_eng.shards = SHARDS;
|
||||
if (wo_engine_start(&mod, heap_mb << 20, nshards) != 0) {
|
||||
fprintf(stderr, "wovm: cannot start %u shards\n", nshards);
|
||||
wo_engine_stop();
|
||||
if (VM.rt.wal) wo_wal_close(&WAL);
|
||||
wo_db_destroy(&DB);
|
||||
wo_vm_destroy(&VM);
|
||||
wo_module_free(&mod);
|
||||
return 2;
|
||||
}
|
||||
}
|
||||
|
||||
/* Program mode: an entry that declares one parameter gets the program's
|
||||
* OWN arguments as a `multi Text` — not the program name, and not the
|
||||
|
|
|
|||
|
|
@ -72,30 +72,43 @@ typedef struct {
|
|||
struct io_uring_sqe *sqes;
|
||||
} rings;
|
||||
|
||||
static struct io_uring_params g_params; /* offsets survive init */
|
||||
/* Ring params live INSIDE each vm (vm.h io_params, opaque bytes) — arc
|
||||
* stage 3 fix: this WAS one shared static ("offsets survive init"), and a
|
||||
* worker's LAZY uring_init memset+refilled it on the worker thread while
|
||||
* another shard was reading ring offsets out of it — submits landed at
|
||||
* garbage offsets, the kernel saw no sqe (enter returned 0, treated as
|
||||
* ok), and a parked fiber's TIMEOUT silently never existed. Per-vm storage
|
||||
* ends the race by construction (a vm's params are only ever touched by
|
||||
* its own thread); the short-submit check below turns any relapse into a
|
||||
* loud trap instead of a lost wake. NOTE: not indexed by shard_id — the
|
||||
* lazy wo_vm_init runs while the worker's shard_id is transiently 0. */
|
||||
_Static_assert(sizeof(struct io_uring_params) <= sizeof(((wo_vm *)0)->io_params),
|
||||
"io_params too small");
|
||||
|
||||
static rings ring_ptrs(const wo_vm *vm) {
|
||||
const struct io_uring_params *p = (const struct io_uring_params *)vm->io_params;
|
||||
rings r;
|
||||
uint8_t *sq = (uint8_t *)vm->io_sq, *cq = (uint8_t *)vm->io_cq;
|
||||
r.sq_head = (uint32_t *)(sq + g_params.sq_off.head);
|
||||
r.sq_tail = (uint32_t *)(sq + g_params.sq_off.tail);
|
||||
r.sq_mask = (uint32_t *)(sq + g_params.sq_off.ring_mask);
|
||||
r.sq_array = (uint32_t *)(sq + g_params.sq_off.array);
|
||||
r.cq_head = (uint32_t *)(cq + g_params.cq_off.head);
|
||||
r.cq_tail = (uint32_t *)(cq + g_params.cq_off.tail);
|
||||
r.cq_mask = (uint32_t *)(cq + g_params.cq_off.ring_mask);
|
||||
r.cqes = (struct io_uring_cqe *)(cq + g_params.cq_off.cqes);
|
||||
r.sq_head = (uint32_t *)(sq + p->sq_off.head);
|
||||
r.sq_tail = (uint32_t *)(sq + p->sq_off.tail);
|
||||
r.sq_mask = (uint32_t *)(sq + p->sq_off.ring_mask);
|
||||
r.sq_array = (uint32_t *)(sq + p->sq_off.array);
|
||||
r.cq_head = (uint32_t *)(cq + p->cq_off.head);
|
||||
r.cq_tail = (uint32_t *)(cq + p->cq_off.tail);
|
||||
r.cq_mask = (uint32_t *)(cq + p->cq_off.ring_mask);
|
||||
r.cqes = (struct io_uring_cqe *)(cq + p->cq_off.cqes);
|
||||
r.sqes = (struct io_uring_sqe *)vm->io_sqes;
|
||||
return r;
|
||||
}
|
||||
|
||||
static int uring_init(wo_vm *vm) {
|
||||
memset(&g_params, 0, sizeof g_params);
|
||||
long fd = syscall(SYS_io_uring_setup, 64u, &g_params);
|
||||
struct io_uring_params *p = (struct io_uring_params *)vm->io_params;
|
||||
memset(p, 0, sizeof *p);
|
||||
long fd = syscall(SYS_io_uring_setup, 64u, p);
|
||||
if (fd < 0) return -1;
|
||||
size_t sq_len = g_params.sq_off.array + g_params.sq_entries * sizeof(uint32_t);
|
||||
size_t cq_len = g_params.cq_off.cqes + g_params.cq_entries * sizeof(struct io_uring_cqe);
|
||||
size_t sqes_len = g_params.sq_entries * sizeof(struct io_uring_sqe);
|
||||
size_t sq_len = p->sq_off.array + p->sq_entries * sizeof(uint32_t);
|
||||
size_t cq_len = p->cq_off.cqes + p->cq_entries * sizeof(struct io_uring_cqe);
|
||||
size_t sqes_len = p->sq_entries * sizeof(struct io_uring_sqe);
|
||||
void *sq = mmap(NULL, sq_len, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_POPULATE, (int)fd,
|
||||
IORING_OFF_SQ_RING);
|
||||
void *cq = mmap(NULL, cq_len, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_POPULATE, (int)fd,
|
||||
|
|
@ -128,7 +141,9 @@ static int uring_submit(wo_vm *vm, const struct io_uring_sqe *sqe) {
|
|||
r.sq_array[idx] = idx;
|
||||
__atomic_store_n(r.sq_tail, tail + 1, __ATOMIC_RELEASE);
|
||||
long rc = syscall(SYS_io_uring_enter, vm->io_fd, 1u, 0u, 0u, NULL, 0);
|
||||
return rc < 0 ? -1 : 0;
|
||||
/* a short submit is a LOST WAKE, never a success (the g_params race
|
||||
* above hid behind rc >= 0 for a whole debugging session) */
|
||||
return rc == 1 ? 0 : -1;
|
||||
}
|
||||
|
||||
/* ---- backend-neutral helpers ------------------------------------------ */
|
||||
|
|
@ -160,6 +175,10 @@ static void wake(wo_vm *vm, wo_fiber *fb) {
|
|||
vm->qtail = fb;
|
||||
}
|
||||
|
||||
void wo_io_unpark(wo_vm *vm, wo_fiber *fb) {
|
||||
if (fb->state == WO_FIB_PARKED) wake(vm, fb);
|
||||
}
|
||||
|
||||
/* ---- API --------------------------------------------------------------- */
|
||||
|
||||
int wo_io_init(wo_vm *vm) {
|
||||
|
|
@ -192,6 +211,9 @@ int wo_io_arm(wo_vm *vm, wo_fiber *fb) {
|
|||
fb->pnext = vm->parked;
|
||||
vm->parked = fb;
|
||||
vm->nparked++;
|
||||
/* arc stage 3: an inbox-wait fiber holds no plane wait at all — the
|
||||
* wake is wo_io_unpark from the envelope drain */
|
||||
if (fb->park_fd == WO_PARK_INBOX) return 0;
|
||||
if (vm->io_kind == 0) {
|
||||
struct io_uring_sqe sqe;
|
||||
memset(&sqe, 0, sizeof sqe);
|
||||
|
|
@ -305,7 +327,7 @@ int wo_io_wait(wo_vm *vm) {
|
|||
int timeout = -1;
|
||||
int64_t now = now_ms();
|
||||
for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext)
|
||||
if (fb->park_fd < 0) {
|
||||
if (fb->park_fd == -1) { /* deadline waits only, never INBOX */
|
||||
int64_t rel = fb->park_deadline - now;
|
||||
if (rel < 0) rel = 0;
|
||||
if (timeout < 0 || rel < timeout) timeout = (int)rel;
|
||||
|
|
@ -333,7 +355,7 @@ int wo_io_wait(wo_vm *vm) {
|
|||
wo_fiber *fb = vm->parked;
|
||||
while (fb) {
|
||||
wo_fiber *nx = fb->pnext;
|
||||
if (fb->park_fd < 0 && fb->park_deadline <= now) {
|
||||
if (fb->park_fd == -1 && fb->park_deadline <= now) {
|
||||
wake(vm, fb);
|
||||
woke = 1;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,9 +22,15 @@ int wo_io_init(wo_vm *vm);
|
|||
void wo_io_destroy(wo_vm *vm);
|
||||
|
||||
/* Register the just-parked fiber's wait (fb->park_* already filled by the
|
||||
* builtin). 0 ok; -1 = arming failed (caller traps the builtin as IO). */
|
||||
* builtin). 0 ok; -1 = arming failed (caller traps the builtin as IO).
|
||||
* park_fd == WO_PARK_INBOX joins the parked list with NO plane wait — the
|
||||
* wake arrives as an inbox envelope (arc stage 3's DB reply). */
|
||||
int wo_io_arm(wo_vm *vm, wo_fiber *fb);
|
||||
|
||||
/* Wake one parked fiber from OUTSIDE the plane (the inbox path: the DB
|
||||
* actor's reply). parked list -> run queue. */
|
||||
void wo_io_unpark(wo_vm *vm, wo_fiber *fb);
|
||||
|
||||
/* Block until at least one parked fiber wakes; woken fibers move to the
|
||||
* run queue. 0 = something woke; WO_IO_STOP = the stop flag interrupted
|
||||
* the wait (caller unwinds everything); -1 = fatal backend error. */
|
||||
|
|
|
|||
220
runtime/src/vm.c
220
runtime/src/vm.c
|
|
@ -12,6 +12,11 @@
|
|||
#include "gc.h"
|
||||
#include "park.h"
|
||||
|
||||
#include <assert.h>
|
||||
|
||||
#include "db.h" /* arc stage 3: the transparent DB RPC (wo_db_req) */
|
||||
#include "table.h" /* slot encode/decode for the RPC marshaling */
|
||||
|
||||
#include <pthread.h>
|
||||
#include <poll.h>
|
||||
#include <sched.h>
|
||||
|
|
@ -94,6 +99,28 @@ static int wo_vm_adopt(wo_vm *vm) {
|
|||
case 2: /* a home-routed free: this arena owns the object */
|
||||
wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->payload);
|
||||
break;
|
||||
case 3: { /* arc stage 3: a marshaled DB statement — WE are the DB
|
||||
* actor (only shard 0 ever receives these). Execute
|
||||
* serialized, right here on the owner thread, then ship
|
||||
* the same request back as the reply. */
|
||||
wo_db_req *q = (wo_db_req *)(uintptr_t)e->payload;
|
||||
assert(vm->is_primary && "DB requests route to shard 0 only");
|
||||
wo_db_exec_req(vm, q);
|
||||
q->done = 1;
|
||||
wo_envelope *re = calloc(1, sizeof *re);
|
||||
if (re) {
|
||||
re->kind = 4;
|
||||
re->payload = e->payload;
|
||||
inbox_push_to(q->from_shard, re);
|
||||
} /* OOM: the requester stays parked until stop — leak, not UB */
|
||||
break;
|
||||
}
|
||||
case 4: { /* the DB actor's reply: wake the requesting fiber; the
|
||||
* re-executed builtin consumes the request */
|
||||
wo_db_req *q = (wo_db_req *)(uintptr_t)e->payload;
|
||||
wo_io_unpark(vm, (wo_fiber *)q->fiber);
|
||||
break;
|
||||
}
|
||||
}
|
||||
free(e);
|
||||
n++;
|
||||
|
|
@ -113,6 +140,190 @@ void wo_route_free(wo_hdr *h) {
|
|||
inbox_push_to(h->shard_id, e);
|
||||
}
|
||||
|
||||
/* ---- arc stage 3: the requester half of the transparent DB RPC ---------
|
||||
* A worker shard's DB builtin lands here (its rt.db is NULL by design):
|
||||
* the args are ENCODED into engine slots on THIS thread — VM heaps are
|
||||
* never read cross-shard — the request rides an envelope to shard 0, and
|
||||
* the fiber parks with no plane wait (WO_PARK_INBOX). The reply unparks
|
||||
* the fiber, the builtin RE-EXECUTES, lands here again, and consumes the
|
||||
* answer. Every status/msg pair is the one the local path would trap. */
|
||||
int wo_db_rpc(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
||||
uint32_t A = wo_ins_a(ins), B = wo_ins_b(ins), C = wo_ins_c(ins);
|
||||
wo_fiber *fb = vm->cur;
|
||||
wo_db_req *q = (wo_db_req *)fb->dbreq;
|
||||
const wo_classdesc *classes = vm->mod->classes;
|
||||
|
||||
if (q && q->done) { /* the reply: consume it and finish the builtin */
|
||||
fb->dbreq = NULL;
|
||||
int rc = q->status;
|
||||
if (rc) {
|
||||
*msg = q->msg;
|
||||
} else {
|
||||
switch (C) {
|
||||
case WO_B_DB_INSERT: R[A] = q->result; break;
|
||||
case WO_B_DB_UPDATE_FIELD:
|
||||
case WO_B_DB_DELETE: R[A] = 0; break;
|
||||
case WO_B_DB_SCAN:
|
||||
case WO_B_DB_PROBE: {
|
||||
wo_multi *ids = wo_multi_new(&vm->rt, WO_K_SCALAR);
|
||||
if (!ids) rc = WO_T_OOM;
|
||||
for (uint32_t i = 0; !rc && i < q->id_cnt; i++)
|
||||
if (wo_multi_push(ids, q->ids[i]) != 0) rc = WO_T_OOM;
|
||||
if (!rc) R[A] = (uint64_t)(uintptr_t)ids;
|
||||
else *msg = "out of memory";
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_GET_FIELD: {
|
||||
int ok = 1;
|
||||
uint64_t v = wo_val_decode_vm(NULL, &vm->rt, q->val_kind, q->val, &ok, msg);
|
||||
wo_db_val_free(NULL, q->val_kind, q->val);
|
||||
q->val = 0;
|
||||
if (!ok) rc = WO_T_OOM;
|
||||
else R[A] = v;
|
||||
break;
|
||||
}
|
||||
default:
|
||||
*msg = "unknown db builtin";
|
||||
rc = WO_T_DB;
|
||||
}
|
||||
}
|
||||
if (q->val) wo_db_val_free(NULL, q->val_kind, q->val);
|
||||
free(q->ids);
|
||||
free(q);
|
||||
return rc;
|
||||
}
|
||||
|
||||
/* first entry: marshal on OUR thread, ship, park */
|
||||
q = calloc(1, sizeof *q);
|
||||
if (!q) {
|
||||
*msg = "out of memory";
|
||||
return WO_T_OOM;
|
||||
}
|
||||
q->op = C;
|
||||
q->from_shard = vm->shard_id;
|
||||
q->fiber = fb;
|
||||
int ok = 1;
|
||||
switch (C) {
|
||||
case WO_B_DB_INSERT: {
|
||||
q->cid = (uint32_t)R[B];
|
||||
if (q->cid >= vm->mod->class_cnt) {
|
||||
free(q);
|
||||
*msg = "no such class";
|
||||
return WO_T_DB;
|
||||
}
|
||||
const wo_classdesc *c = &classes[q->cid];
|
||||
q->slots = calloc(c->field_cnt ? c->field_cnt : 1, 8);
|
||||
if (!q->slots) {
|
||||
free(q);
|
||||
*msg = "out of memory";
|
||||
return WO_T_OOM;
|
||||
}
|
||||
q->slot_cnt = c->field_cnt;
|
||||
for (uint32_t i = 0; i < c->field_cnt; i++) {
|
||||
q->slots[i] = wo_db_val_encode(classes, c->kinds[i], R[B + 1 + i], &ok, msg);
|
||||
if (!ok) { /* GCREF (the compiler's reject, defensively) or OOM */
|
||||
for (uint32_t j = 0; j < i; j++)
|
||||
wo_db_val_free(NULL, c->kinds[j], q->slots[j]);
|
||||
free(q->slots);
|
||||
free(q);
|
||||
return WO_T_DB;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_UPDATE_FIELD: {
|
||||
q->cid = (uint32_t)R[B];
|
||||
q->id = R[B + 1];
|
||||
q->field = (uint32_t)R[B + 2];
|
||||
if (q->cid >= vm->mod->class_cnt || q->field >= classes[q->cid].field_cnt) {
|
||||
free(q);
|
||||
*msg = "no such field";
|
||||
return WO_T_DB;
|
||||
}
|
||||
q->slots = calloc(1, 8);
|
||||
if (!q->slots) {
|
||||
free(q);
|
||||
*msg = "out of memory";
|
||||
return WO_T_OOM;
|
||||
}
|
||||
q->slot_cnt = 1;
|
||||
q->slots[0] =
|
||||
wo_db_val_encode(classes, classes[q->cid].kinds[q->field], R[B + 3], &ok, msg);
|
||||
if (!ok) {
|
||||
free(q->slots);
|
||||
free(q);
|
||||
return WO_T_DB;
|
||||
}
|
||||
break;
|
||||
}
|
||||
case WO_B_DB_DELETE:
|
||||
q->cid = (uint32_t)R[B];
|
||||
q->id = R[B + 1];
|
||||
break;
|
||||
case WO_B_DB_SCAN:
|
||||
q->cid = (uint32_t)R[B];
|
||||
break;
|
||||
case WO_B_DB_GET_FIELD:
|
||||
q->cid = (uint32_t)R[B];
|
||||
q->id = R[B + 1];
|
||||
q->field = (uint32_t)R[B + 2];
|
||||
break;
|
||||
case WO_B_DB_PROBE: {
|
||||
q->cid = (uint32_t)R[B];
|
||||
q->index = (uint32_t)R[B + 1];
|
||||
/* the key's kind comes from the class table's index metadata —
|
||||
* identical on every shard (one module). An index the metadata
|
||||
* does not know ships keyless; the owner answers empty, exactly
|
||||
* as the local path does. */
|
||||
if (q->cid < vm->mod->class_cnt && classes[q->cid].idx_meta &&
|
||||
q->index < classes[q->cid].idx_cnt) {
|
||||
const uint32_t *p = classes[q->cid].idx_meta;
|
||||
for (uint32_t k = 0; k < q->index; k++) p += 2 + p[1];
|
||||
uint8_t kind = classes[q->cid].kinds[p[2]];
|
||||
q->slots = calloc(1, 8);
|
||||
if (!q->slots) {
|
||||
free(q);
|
||||
*msg = "out of memory";
|
||||
return WO_T_OOM;
|
||||
}
|
||||
q->slot_cnt = 1;
|
||||
q->slots[0] = wo_db_val_encode(classes, kind, R[B + 2], &ok, msg);
|
||||
if (!ok) {
|
||||
free(q->slots);
|
||||
free(q);
|
||||
return WO_T_DB;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
default:
|
||||
free(q);
|
||||
*msg = "unknown db builtin";
|
||||
return WO_T_DB;
|
||||
}
|
||||
wo_envelope *e = calloc(1, sizeof *e);
|
||||
if (!e) {
|
||||
if (q->slot_cnt && C == WO_B_DB_INSERT) {
|
||||
const wo_classdesc *c = &classes[q->cid];
|
||||
for (uint32_t j = 0; j < c->field_cnt; j++)
|
||||
wo_db_val_free(NULL, c->kinds[j], q->slots[j]);
|
||||
} else if (q->slot_cnt && C == WO_B_DB_UPDATE_FIELD) {
|
||||
wo_db_val_free(NULL, classes[q->cid].kinds[q->field], q->slots[0]);
|
||||
} /* a PROBE key leaks on this path: kind recompute not worth it */
|
||||
free(q->slots);
|
||||
free(q);
|
||||
*msg = "out of memory";
|
||||
return WO_T_OOM;
|
||||
}
|
||||
fb->dbreq = q;
|
||||
e->kind = 3;
|
||||
e->payload = (uint64_t)(uintptr_t)q;
|
||||
inbox_push_to(0, e);
|
||||
fb->park_fd = WO_PARK_INBOX;
|
||||
fb->park_done = 0; /* resume RE-EXECUTES the builtin: the consume path */
|
||||
return WO_SYS_PARKED;
|
||||
}
|
||||
|
||||
/* A worker's whole life in T5: pinned, parked on its wake eventfd until
|
||||
* shutdown. T6 gives it an inbox to adopt fibers from and the serve loop
|
||||
* that runs them. */
|
||||
|
|
@ -821,6 +1032,11 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) {
|
|||
do { \
|
||||
if (--vm->budget <= 0) { \
|
||||
vm->budget = vm->budget0; \
|
||||
/* arc stage 3: a busy shard still serves its inbox once per \
|
||||
* slice — bounds a DB request's wait on a computing primary \
|
||||
* to one reduction budget */ \
|
||||
if (INBOX_READY[vm->shard_id % WO_ENG_MAX_SHARDS]) \
|
||||
(void)wo_vm_adopt(vm); \
|
||||
if (vm->qhead) { \
|
||||
vm->cur->frames[vm->cur->depth - 1].pc = pc; \
|
||||
fib_enqueue(vm, vm->cur); \
|
||||
|
|
@ -1349,6 +1565,10 @@ dispatch:
|
|||
* stopped (1), or a fatal error (-1). */
|
||||
int wo_vm_serve(wo_vm *vm) {
|
||||
tls_vm = vm;
|
||||
/* arc stage 3 obligation: a worker NEVER holds the engine or the WAL —
|
||||
* its DB statements marshal to shard 0 (wo_db_rpc). Replay finished on
|
||||
* the primary before wo_engine_start spawned this thread. */
|
||||
assert(!vm->rt.db && !vm->rt.wal);
|
||||
if (!vm->qhead) return 2;
|
||||
vm->cur = fib_dequeue(vm);
|
||||
vm->budget = vm->budget0;
|
||||
|
|
|
|||
|
|
@ -79,8 +79,16 @@ typedef struct wo_fiber {
|
|||
* borrows — the RUNTIME owns it and drops it after the call returns. */
|
||||
struct wo_actor *actor;
|
||||
uint64_t cur_msg;
|
||||
/* arc stage 3: the in-flight DB request while parked on the DB actor's
|
||||
* reply (a wo_db_req*, opaque here; vm.c owns the protocol) */
|
||||
void *dbreq;
|
||||
} wo_fiber;
|
||||
|
||||
/* arc stage 3: park_fd sentinel — PARKED with NO plane wait; the wake is
|
||||
* an inbox envelope (the DB actor's reply). Excluded from the deadline
|
||||
* scans, which key on park_fd == -1 exactly. */
|
||||
#define WO_PARK_INBOX (-2)
|
||||
|
||||
/* An actor: moved-in state, its receive method, a FIFO mailbox, and at
|
||||
* most one delivery fiber at a time (one message at a time — the actor
|
||||
* guarantee). Actors live until program end (v1: no actor death). */
|
||||
|
|
@ -126,6 +134,12 @@ typedef struct wo_vm {
|
|||
int io_fd; /* ring fd or epoll fd */
|
||||
void *io_sq, *io_cq, *io_sqes; /* uring mmaps (NULL under epoll) */
|
||||
size_t io_sq_len, io_cq_len, io_sqes_len;
|
||||
/* THIS ring's io_uring_params (opaque bytes; park.c owns the type).
|
||||
* Arc stage 3 fix: a single file-static params was rewritten by every
|
||||
* shard's lazy init while other shards read ring offsets out of it —
|
||||
* submits landed at garbage offsets and parked fibers lost their
|
||||
* wakes. Per-vm storage ends the race by construction. */
|
||||
unsigned char io_params[256];
|
||||
} wo_vm;
|
||||
|
||||
/* arc: the spawn/send builtins' runtime halves (vm.c owns the scheduler). */
|
||||
|
|
@ -146,14 +160,23 @@ typedef struct wo_engine {
|
|||
|
||||
extern wo_engine wo_eng; /* the process's one engine (vm.c) */
|
||||
|
||||
/* inbox envelope kinds (arc T6) */
|
||||
/* inbox envelope kinds (arc T6; 3/4 = stage 3's transparent DB RPC) */
|
||||
typedef struct wo_envelope {
|
||||
struct wo_envelope *next;
|
||||
int kind; /* 0 = SEND (actor, payload), 1 = SPAWN-ADOPT (actor), 2 = FREE (payload = wo_hdr*) */
|
||||
int kind; /* 0 = SEND (actor, payload), 1 = SPAWN-ADOPT (actor),
|
||||
* 2 = FREE (payload = wo_hdr*),
|
||||
* 3 = DB_REQ (payload = wo_db_req*, to shard 0),
|
||||
* 4 = DB_RESP (payload = wo_db_req*, back to the requester) */
|
||||
struct wo_actor *actor;
|
||||
uint64_t payload;
|
||||
} wo_envelope;
|
||||
|
||||
/* arc stage 3: the requester half of the transparent DB RPC (vm.c). Called
|
||||
* by the builtin dispatcher on a worker shard whose rt.db is NULL: first
|
||||
* entry marshals + parks (WO_SYS_PARKED), the re-execution after the reply
|
||||
* consumes it. */
|
||||
int wo_db_rpc(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg);
|
||||
|
||||
/* the shard whose thread we are on (thread-local; obj.c stamps and gc.c
|
||||
* routes with it). NULL only before main's vm exists. */
|
||||
wo_vm *wo_tls_vm(void);
|
||||
|
|
|
|||
72
scripts/db-actor-accept.sh
Executable file
72
scripts/db-actor-accept.sh
Executable file
|
|
@ -0,0 +1,72 @@
|
|||
#!/usr/bin/env bash
|
||||
# scripts/db-actor-accept.sh — arc stage 3's gate: actors on worker shards
|
||||
# read and write the database through the transparent DB actor. Multi-shard
|
||||
# output is asserted as a SET (scheduling orders the writer lines); the
|
||||
# main line and the single-shard run are exact. The WO_DATA pair proves a
|
||||
# worker's write rides the owner's WAL and replays.
|
||||
set -uo pipefail
|
||||
|
||||
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
WOC="$ROOT/compiler/_build/default/bin/woc"
|
||||
WOVM="$ROOT/runtime/wovm"
|
||||
DIR="$ROOT/docs/examples/db-actor"
|
||||
|
||||
pass=0; fail=0
|
||||
ok() { echo "ok $1"; pass=$((pass + 1)); }
|
||||
bad() { echo "FAIL $1 -- $2"; fail=$((fail + 1)); }
|
||||
|
||||
if [[ ! -x "$WOC" || ! -x "$WOVM" ]]; then
|
||||
echo "db-actor-accept: build woc and wovm first" >&2; exit 1
|
||||
fi
|
||||
|
||||
if "$WOC" build "$DIR" -o "$DIR/target/db-actor" --runtime "$WOVM" >/dev/null 2>&1; then
|
||||
ok "builds"
|
||||
else
|
||||
bad "build" "woc failed"; echo "db-actor-accept: 1 checks, 1 failures"; exit 1
|
||||
fi
|
||||
|
||||
check_set() { # name [env pairs...]
|
||||
local name="$1"; shift
|
||||
local out
|
||||
out="$(env "$@" timeout 30 "$DIR/target/db-actor" 2>&1)"
|
||||
if printf '%s' "$out" | grep -q "writer 1 sees sum" \
|
||||
&& printf '%s' "$out" | grep -q "writer 2 sees sum" \
|
||||
&& printf '%s' "$out" | grep -q "^main sees 2 rows, sum 3$" ; then
|
||||
ok "$name: both writers wrote and read cross-shard; main exact"
|
||||
else
|
||||
bad "$name" "$(printf '%s' "$out" | tr '\n' '|')"
|
||||
fi
|
||||
}
|
||||
|
||||
# multi-shard (default = all cores): the RPC path under test, three rounds
|
||||
check_set "multi #1"
|
||||
check_set "multi #2"
|
||||
check_set "multi #3"
|
||||
# forced backends: the reply park is plane-independent
|
||||
check_set "multi uring" WO_IO=uring
|
||||
check_set "multi epoll" WO_IO=epoll
|
||||
|
||||
# single-shard: byte-exact — the local path is untouched
|
||||
sout="$(WO_SHARDS=1 timeout 30 "$DIR/target/db-actor" 2>&1)"
|
||||
want=$'writer 1 sees sum 1\nwriter 2 sees sum 3\nmain sees 2 rows, sum 3'
|
||||
if [[ "$sout" == "$want" ]]; then
|
||||
ok "single-shard byte-exact"
|
||||
else
|
||||
bad "single-shard" "$(printf '%s' "$sout" | tr '\n' '|')"
|
||||
fi
|
||||
|
||||
# durability through the RPC: a worker's insert commits on the owner's WAL
|
||||
# before the ack; a restart replays it (2 rows, then 2+2)
|
||||
DATA="$(mktemp -d)"
|
||||
r1="$(WO_DATA="$DATA" timeout 30 "$DIR/target/db-actor" 2>&1 | tail -1)"
|
||||
r2="$(WO_DATA="$DATA" timeout 30 "$DIR/target/db-actor" 2>&1 | tail -1)"
|
||||
rm -rf "$DATA"
|
||||
if [[ "$r1" == "main sees 2 rows, sum 3" && "$r2" == "main sees 4 rows, sum 6" ]]; then
|
||||
ok "WAL: worker writes ack-after-durable, replay doubles the store"
|
||||
else
|
||||
bad "WAL replay" "r1=$r1 r2=$r2"
|
||||
fi
|
||||
|
||||
echo
|
||||
echo "db-actor-accept: $((pass + fail)) checks, $fail failures"
|
||||
[[ $fail -eq 0 ]] || exit 1
|
||||
Loading…
Reference in a new issue