From de714f5023a68799e66a4b08711ee1992c9b18af Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 27 Sep 2026 16:17:04 +0300 Subject: [PATCH] Add bounded Merkle pull lifecycle and terminal error handoff --- src/chat/merkle_sync.c | 82 ++++++++++++++++++++++++++++++++---- src/chat/merkle_sync.h | 9 ++++ tests/test_merkle_protocol.c | 70 ++++++++++++++++++++++++++++++ 3 files changed, 152 insertions(+), 9 deletions(-) diff --git a/src/chat/merkle_sync.c b/src/chat/merkle_sync.c index 08048476..5f1b3376 100644 --- a/src/chat/merkle_sync.c +++ b/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) { diff --git a/src/chat/merkle_sync.h b/src/chat/merkle_sync.h index 64f7c546..76ed7691 100644 --- a/src/chat/merkle_sync.h +++ b/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); diff --git a/tests/test_merkle_protocol.c b/tests/test_merkle_protocol.c index fbf2d036..4329b71a 100644 --- a/tests/test_merkle_protocol.c +++ b/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;