Browse Source

Add bounded Merkle pull lifecycle and terminal error handoff

proxy
evgeny 4 days ago
parent
commit
de714f5023
  1. 82
      src/chat/merkle_sync.c
  2. 9
      src/chat/merkle_sync.h
  3. 70
      tests/test_merkle_protocol.c

82
src/chat/merkle_sync.c

@ -16,8 +16,9 @@
* Меньший node_id ведёт раунд: pull -> TURN -> reverse pull -> CHECK -> DONE. */
#define MS_HEADER 22
#define MS_QUERY_SIZE 18
#define MS_PROGRESS_TIMEOUT_TB (30u * 10000u)
enum { MS_WAKE = 0x10, MS_BEGIN, MS_READY, MS_QUERY, MS_HASHES, MS_PAGE, MS_TURN, MS_CHECK, MS_DONE, MS_ERROR };
enum { MS_IDLE, MS_BEGIN_WAIT, MS_PULL, MS_SERVE, MS_DONE_WAIT, MS_FAILED };
enum { MS_IDLE, MS_BEGIN_WAIT, MS_PULL, MS_SERVE, MS_DONE_WAIT, MS_FAILED, MS_FAILING };
struct ms_request {
struct ms_request* next;
@ -39,6 +40,9 @@ struct ms_session {
struct queue_waiter_handle waiter;
struct ll_entry* tx;
void* work;
void* deadline;
void* pull;
int failure_error;
struct ms_request* requests;
uint64_t ns_id, peer, round, revision, verified_revision;
char ns[21];
@ -103,6 +107,32 @@ static int ms_schedule(struct ms_session* s);
static void ms_work(void* arg);
static void ms_fail(struct ms_session* s, int error, const char* reason, int report);
static void ms_abort_pull(struct ms_session* s) {
if (!s->pull) return;
void* pull = s->pull; s->pull = NULL;
s->ms->ops->abort_pull(s->ms->data_ctx, pull);
}
static void ms_stop_deadline(struct ms_session* s) {
if (s->deadline) uasync_cancel_timeout(s->ms->inst->ua, s->deadline);
s->deadline = NULL;
}
static void ms_timeout(void* arg) {
struct ms_session* s = arg;
s->deadline = NULL;
ms_fail(s, MT_ERR_TIMEOUT, "no protocol progress for 30 seconds", 0);
}
static int ms_progress(struct ms_session* s) {
ms_stop_deadline(s);
s->deadline = uasync_set_timeout(s->ms->inst->ua, MS_PROGRESS_TIMEOUT_TB, s, ms_timeout, "merkle_progress");
if (!s->deadline) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: deadline allocation ns=%s", s->ns); return MT_ERR_IO;
}
return 0;
}
static void ms_discard_tx(struct ms_session* s) {
queue_waiter_cancel(s->send_q, &s->waiter);
if (s->tx) { queue_dgram_free(s->tx); queue_entry_free(s->tx); s->tx = NULL; }
@ -114,6 +144,8 @@ static void ms_free(struct ms_session* s) {
while (*p && *p != s) p = &(*p)->next;
if (*p) *p = s->next;
if (s->work) uasync_call_soon_cancel(ms->inst->ua, s->work);
ms_stop_deadline(s);
ms_abort_pull(s);
ms_discard_tx(s);
while (s->requests) {
struct ms_request* r = s->requests;
@ -166,8 +198,15 @@ static void ms_sent(struct ll_queue* q, void* arg) {
s->sent_bytes += len;
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: TX ns=%s peer=%016llx round=%llu type=%u bytes=%zu",
s->ns, (unsigned long long)s->peer, (unsigned long long)s->round, type, len);
if (type == MS_ERROR) {
s->state = MS_FAILED;
ms_stop_deadline(s);
ms_notify(s, s->failure_error);
return;
}
if (type == MS_DONE) {
s->state = MS_IDLE;
ms_stop_deadline(s);
if (s->revision == s->verified_revision) {
s->dirty = 0;
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: converged ns=%s peer=%016llx round=%llu tx=%llu rx=%llu",
@ -197,19 +236,24 @@ static int ms_send(struct ms_session* s, int type, uint32_t request, const uint8
s->tx = e;
/* Transport normalizer configures deferred waiters. Не вызываем callback рекурсивно. */
if (queue_waiter_wait(s->send_q, &s->waiter, ms_sent, s) < 0) { ms_discard_tx(s); return -1; }
if (ms_progress(s) < 0) { ms_discard_tx(s); return -1; }
return 0;
}
static void ms_fail(struct ms_session* s, int error, const char* reason, int report) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: failed ns=%s peer=%016llx round=%llu error=%d reason=%s",
s->ns, (unsigned long long)s->peer, (unsigned long long)s->round, error, reason);
s->state = MS_FAILED; s->dirty = 0; s->waiting = 0;
s->state = MS_FAILED; s->dirty = 0; s->waiting = 0; s->failure_error = error;
if (s->work) { uasync_call_soon_cancel(s->ms->inst->ua, s->work); s->work = NULL; }
ms_discard_tx(s);
ms_stop_deadline(s);
ms_abort_pull(s);
if (report) {
s->state = MS_FAILING;
uint8_t code[4]; ms_write32(code, (uint32_t)(-error));
if (ms_send(s, MS_ERROR, 0, code, sizeof(code)) < 0)
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: could not send error ns=%s", s->ns);
if (ms_send(s, MS_ERROR, 0, code, sizeof(code)) == 0) return;
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: could not send error ns=%s", s->ns);
s->state = MS_FAILED;
}
ms_notify(s, error);
}
@ -234,6 +278,9 @@ static int ms_query(struct ms_session* s) {
}
static int ms_pull_start(struct ms_session* s) {
ms_abort_pull(s);
if (s->ms->ops->begin_pull && s->ms->ops->begin_pull(s->ms->data_ctx, s->ns, s->peer, s->round, &s->pull) < 0)
return MT_ERR_DATA;
s->state = MS_PULL; s->depth = 0; s->after_valid = 0; s->after = 0;
memset(s->stack, 0, sizeof(s->stack));
return ms_query(s);
@ -254,6 +301,11 @@ static int ms_walk_next(struct ms_session* s) {
s->depth--;
}
s->waiting = 0;
if (s->pull) {
int rc = s->ms->ops->commit_pull(s->ms->data_ctx, s->pull);
ms_abort_pull(s);
if (rc < 0) return rc;
}
if (s->leader) { s->state = MS_SERVE; return ms_send(s, MS_TURN, 0, NULL, 0); }
if (s->ms->ops->finish && s->ms->ops->finish(s->ms->data_ctx, s->ns) < 0) return MT_ERR_DATA;
uint8_t root[MT_HASH_SIZE];
@ -263,6 +315,7 @@ static int ms_walk_next(struct ms_session* s) {
}
static int ms_begin(struct ms_session* s) {
ms_abort_pull(s);
s->round = ++s->ms->round_seq;
if (!s->round) return -1;
s->request_id = s->remote_request_id = 0;
@ -362,7 +415,8 @@ static int ms_receive_page(struct ms_session* s, uint32_t request, const uint8_t
uint64_t next = ms_read64(p + 1);
if (f->level != MT_MAX_LEVEL || (p[0] && (len == 9 || merkle_sync_level_prefix(next, MT_MAX_LEVEL) != f->prefix
|| (s->after_valid && next <= s->after)))) return MT_ERR_PROTOCOL;
int rc = s->ms->ops->apply_items(s->ms->data_ctx, s->ns, s->peer, p + 9, len - 9);
int rc = s->pull ? s->ms->ops->stage_page(s->ms->data_ctx, s->pull, f->prefix, s->after_valid, s->after,
next, p[0], p + 9, len - 9) : s->ms->ops->apply_items(s->ms->data_ctx, s->ns, s->peer, p + 9, len - 9);
if (rc < 0) return rc == MT_ERR_CONFLICT ? rc : MT_ERR_DATA;
s->waiting = 0;
if (p[0]) { s->after = next; s->after_valid = 1; return ms_query(s) < 0 ? MT_ERR_IO : 0; }
@ -397,6 +451,7 @@ static void ms_receive(struct ETCP_CONN* conn, struct ll_entry* entry) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: RX ns=%s peer=%016llx round=%llu type=%u request=%u bytes=%zu state=%d",
ns, (unsigned long long)s->peer, (unsigned long long)round, type, request, len, s->state);
int rc = 0, notify = 0;
if (s->state == MS_FAILING) goto done;
if (type != MS_WAKE && type != MS_BEGIN && round != s->round) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: discard old round=%llu current=%llu", (unsigned long long)round,
(unsigned long long)s->round);
@ -416,19 +471,20 @@ static void ms_receive(struct ETCP_CONN* conn, struct ll_entry* entry) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: discard duplicate BEGIN ns=%s", ns);
} else {
ms_discard_tx(s);
ms_abort_pull(s);
s->round = round; s->state = MS_SERVE; s->waiting = 0;
s->request_id = s->remote_request_id = 0; s->dirty = 0;
rc = ms_send(s, MS_READY, 0, NULL, 0) < 0 ? MT_ERR_IO : 0;
}
} else if (type == MS_READY) {
if (!s->leader || s->state != MS_BEGIN_WAIT || request || len) rc = MT_ERR_PROTOCOL;
else rc = ms_pull_start(s) < 0 ? MT_ERR_IO : 0;
else rc = ms_pull_start(s);
} else if (type == MS_QUERY) rc = ms_serve_query(s, request, p, len);
else if (type == MS_HASHES) rc = ms_receive_hashes(s, request, p, len);
else if (type == MS_PAGE) rc = ms_receive_page(s, request, p, len);
else if (type == MS_TURN) {
if (s->leader || s->state != MS_SERVE || request || len) rc = MT_ERR_PROTOCOL;
else rc = ms_pull_start(s) < 0 ? MT_ERR_IO : 0;
else rc = ms_pull_start(s);
} else if (type == MS_CHECK) {
uint8_t roots[2 * MT_HASH_SIZE];
if (!s->leader || s->state != MS_SERVE || request || len != MT_HASH_SIZE) rc = MT_ERR_PROTOCOL;
@ -453,6 +509,7 @@ static void ms_receive(struct ETCP_CONN* conn, struct ll_entry* entry) {
else if (ms_root(s, root) < 0) rc = MT_ERR_IO;
else {
s->state = MS_IDLE;
ms_stop_deadline(s);
s->dirty = memcmp(root, p, sizeof(root)) != 0;
if (s->dirty) rc = ms_schedule(s);
else {
@ -462,7 +519,7 @@ static void ms_receive(struct ETCP_CONN* conn, struct ll_entry* entry) {
}
}
} else if (type == MS_ERROR) {
if (request || len != 4 || ms_read32(p) < 1 || ms_read32(p) > 5) rc = MT_ERR_PROTOCOL;
if (request || len != 4 || ms_read32(p) < 1 || ms_read32(p) > 6) rc = MT_ERR_PROTOCOL;
else {
rc = -(int)ms_read32(p);
queue_dgram_free(entry); queue_entry_free(entry);
@ -499,7 +556,10 @@ static void ms_conn_status(struct ETCP_CONN* conn, int event, void* arg) {
}
int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, const struct merkle_sync_data_ops* ops, void* data_ctx) {
if (!inst || inst->msync || !inst->ua || !inst->topo_sqlite_db || !ops || !ops->update_bucket_hash || !ops->get_page || !ops->apply_items) {
if (!inst || inst->msync || !inst->ua || !inst->topo_sqlite_db || !ops || !ops->update_bucket_hash || !ops->get_page
|| (!ops->apply_items && !ops->stage_page)
|| ((ops->begin_pull || ops->commit_pull || ops->abort_pull || ops->stage_page)
&& !(ops->begin_pull && ops->commit_pull && ops->abort_pull && ops->stage_page))) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: invalid initialization"); return -1;
}
if (merkle_tree_init(inst->topo_sqlite_db) < 0) return -1;
@ -542,6 +602,10 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns,
struct ms_session* s = ms_find(ms, conn, ns_id);
if (!s) s = ms_create(ms, conn, ns_id);
if (!s) { u_free(r); return -1; }
if (s->state == MS_FAILING) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_sync: start during terminal handoff ns=%s", ns);
u_free(r); return -1;
}
if (s->state == MS_FAILED) { s->state = MS_IDLE; s->previous_valid = 0; }
s->dirty = 1;
if (s->state == MS_IDLE && ms_schedule(s) < 0) {

9
src/chat/merkle_sync.h

@ -12,6 +12,7 @@ struct UTUN_INSTANCE;
#define MT_ERR_DATA -3
#define MT_ERR_CONFLICT -4
#define MT_ERR_DISCONNECTED -5
#define MT_ERR_TIMEOUT -6
#define MT_PAGE_SIZE 16384
/* Namespace — каноническая десятичная строка uint64_t. Движок не знает CHAT-групп.
@ -30,6 +31,14 @@ struct merkle_sync_data_ops {
int (*validate_peer)(void* ctx, const char* ns, uint64_t peer);
/* Проверить отложенные зависимости модели после обоих проходов, до подтверждения корня. */
int (*finish)(void* ctx, const char* ns);
/* Optional pull transaction. All three hooks must be supplied together.
* abort_pull releases the context after both success and failure.
* Staged records must not affect published hashes before commit_pull. */
int (*begin_pull)(void* ctx, const char* ns, uint64_t peer, uint64_t round, void** pull);
int (*commit_pull)(void* ctx, void* pull);
void (*abort_pull)(void* ctx, void* pull);
int (*stage_page)(void* ctx, void* pull, uint64_t prefix, int after_valid, uint64_t after,
uint64_t next, int more, const uint8_t* data, size_t len);
};
typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, void* arg);

70
tests/test_merkle_protocol.c

@ -89,6 +89,8 @@ static void test_malformed_and_old_round(void) {
t_receive(&f, MS_HASHES, 6, 1, short_hashes, sizeof(short_hashes));
assert(!f.callbacks && s->waiting && s->round == 7);
t_receive(&f, MS_HASHES, 7, 1, short_hashes, sizeof(short_hashes));
assert(!f.callbacks && s->state == MS_FAILING);
uasync_poll(f.inst.ua, 0);
assert(f.callbacks == 1 && f.result == MT_ERR_PROTOCOL && s->state == MS_FAILED);
t_destroy(&f);
puts("PASS: stale round ignored; truncated hashes rejected before reading");
@ -100,6 +102,8 @@ static void test_rejected_data(void) {
f.apply_error = MT_ERR_DATA;
uint8_t page[11] = {0};
t_receive(&f, MS_PAGE, 7, 1, page, sizeof(page));
assert(!f.callbacks && s->state == MS_FAILING);
uasync_poll(f.inst.ua, 0);
assert(f.callbacks == 1 && f.result == MT_ERR_DATA && s->state == MS_FAILED);
t_destroy(&f);
puts("PASS: failed apply cannot produce success");
@ -203,12 +207,78 @@ static void test_send_rejection(void) {
puts("PASS: transport rejection retains ownership; packet is freed exactly once");
}
static void t_cancel_done(uint64_t peer, const char* ns, int result, void* arg) {
struct fixture* f = arg;
t_done(peer, ns, result, arg);
merkle_sync_cancel(&f->inst, peer, ns);
}
static void test_error_handoff(void) {
struct fixture f; t_init(&f);
struct ms_session* s = t_session(&f); s->requests->cb = t_cancel_done;
s->stack[0].level = MT_MAX_LEVEL; f.apply_error = MT_ERR_DATA;
uint8_t page[11] = {0};
t_receive(&f, MS_PAGE, 7, 1, page, sizeof(page));
assert(!f.callbacks && s->state == MS_FAILING);
uasync_poll(f.inst.ua, 0);
assert(f.callbacks == 1 && f.result == MT_ERR_DATA && !f.inst.msync->sessions);
assert(f.conn.send_input_q->head && f.conn.send_input_q->head->dgram[9] == MS_ERROR);
t_destroy(&f);
puts("PASS: ERROR reaches transport before cancellation inside callback");
}
static void test_blocked_error_timeout(void) {
struct fixture f; t_init(&f);
struct ll_entry* blocker = queue_entry_new(0); assert(blocker);
assert(queue_data_put(f.conn.send_input_q, blocker) == 0);
struct ms_session* s = t_session(&f); s->requests->cb = t_cancel_done;
ms_fail(s, MT_ERR_DATA, "injected apply failure", 1);
assert(s->deadline && s->state == MS_FAILING && !f.callbacks);
ms_stop_deadline(s); ms_timeout(s);
assert(f.callbacks == 1 && f.result == MT_ERR_TIMEOUT && !f.inst.msync->sessions);
uasync_poll(f.inst.ua, 0);
assert(f.callbacks == 1 && !f.conn.send_input_q->waiter_head);
t_destroy(&f);
puts("PASS: blocked terminal packet times out once and releases waiter/session");
}
static int pull_created, pull_committed, pull_released;
static int t_begin_pull(void* ctx, const char* ns, uint64_t peer, uint64_t round, void** pull) {
(void)ctx; (void)ns; (void)peer; (void)round;
*pull = u_malloc(1); assert(*pull); pull_created++; return 0;
}
static int t_commit_pull(void* ctx, void* pull) { (void)ctx; assert(pull); pull_committed++; return 0; }
static void t_abort_pull(void* ctx, void* pull) { (void)ctx; assert(pull); pull_released++; u_free(pull); }
static int t_stage_page(void* ctx, void* pull, uint64_t prefix, int av, uint64_t after,
uint64_t next, int more, const uint8_t* data, size_t len) {
(void)ctx; (void)pull; (void)prefix; (void)av; (void)after; (void)next; (void)more; (void)data; (void)len; return 0;
}
static void test_pull_lifecycle(void) {
struct fixture f; t_init(&f);
struct merkle_sync_data_ops ops = t_ops;
ops.begin_pull = t_begin_pull; ops.commit_pull = t_commit_pull; ops.abort_pull = t_abort_pull; ops.stage_page = t_stage_page;
f.inst.msync->ops = &ops;
struct ms_session* s = t_session(&f);
assert(ms_pull_start(s) == 0 && s->pull);
ms_discard_tx(s); s->stack[0].level = MT_MAX_LEVEL;
assert(ms_walk_next(s) == 0 && !s->pull && pull_committed == 1 && pull_released == 1);
ms_discard_tx(s);
assert(ms_pull_start(s) == 0);
ms_discard_tx(s); assert(ms_begin(s) == 0 && !s->pull && pull_released == 2);
ms_discard_tx(s); assert(ms_pull_start(s) == 0);
ms_conn_status(&f.conn, ETCP_CONN_STATUS_DOWN, f.inst.msync);
assert(pull_created == 3 && pull_committed == 1 && pull_released == 3);
t_destroy(&f);
puts("PASS: pull commits once; restart and disconnect abort exactly their own context");
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
size_t baseline = u_get_allocated_count();
test_malformed_and_old_round(); test_rejected_data(); test_sibling_walk(); test_cancel_waiters();
test_page_size_and_cursor(); test_readonly_tree(); test_backpressure_and_disconnect(); test_peer_restart();
test_send_rejection();
test_error_handoff(); test_blocked_error_timeout(); test_pull_lifecycle();
assert(u_get_allocated_count() == baseline);
puts("ALL PASS: protocol regressions, no tracked allocations leaked");
return 0;

Loading…
Cancel
Save