|
|
|
|
@ -48,6 +48,7 @@
|
|
|
|
|
|
|
|
|
|
#define TICK_TB 100 /* 10 ms на такт state machine */ |
|
|
|
|
#define MAX_TICKS 3000 /* глобальный таймаут ≈ 30 c */ |
|
|
|
|
#define OFFLINE_BATCH 40 /* Больше лимита одной порции durable pump. */ |
|
|
|
|
|
|
|
|
|
/* фазы master state-machine */ |
|
|
|
|
enum dm_phase { |
|
|
|
|
@ -232,6 +233,20 @@ static int outbox_count(struct UTUN_INSTANCE* inst) {
|
|
|
|
|
return count; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Все 40 последовательных сообщений сохранены; максимум seq сам по себе этого не доказывает. */ |
|
|
|
|
static int offline_batch_received(struct UTUN_INSTANCE* inst) { |
|
|
|
|
sqlite3_stmt* st = NULL; |
|
|
|
|
int count = -1; |
|
|
|
|
if (sqlite3_prepare_v2(inst->topo_sqlite_db, |
|
|
|
|
"SELECT COUNT(*) FROM dm_messages WHERE conv_id=? AND dir=0 AND seq>=2 AND seq<?", -1, &st, NULL) == SQLITE_OK) { |
|
|
|
|
sqlite3_bind_text(st, 1, g_sh.conv_id, -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_int(st, 2, OFFLINE_BATCH + 2); |
|
|
|
|
if (sqlite3_step(st) == SQLITE_ROW) count = sqlite3_column_int(st, 0); |
|
|
|
|
} |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
return count == OFFLINE_BATCH; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Каноническое сообщение реального автора для проверок порядка, подписи и commit. */ |
|
|
|
|
static size_t test_body(struct UTUN_INSTANCE* author, uint64_t seq, const char* text, uint8_t body[512]) { |
|
|
|
|
uint64_t conv = strtoull(g_sh.conv_id, NULL, 10), ts = 1; |
|
|
|
|
@ -455,21 +470,26 @@ static void dm_tick(void* arg) {
|
|
|
|
|
if (!bgp_has(A, nb)) t->phase = P_SEND2; |
|
|
|
|
break; |
|
|
|
|
case P_SEND2: |
|
|
|
|
if (dm_send(A, g_sh.conv_id, "text", (const uint8_t*)"hello2", 6) != 0) { |
|
|
|
|
fprintf(stderr, "A: dm_send hello2 failed\n"); |
|
|
|
|
t->result = 2; |
|
|
|
|
} else { |
|
|
|
|
t->phase = P_WAIT_MAIL; |
|
|
|
|
for (int i = 0; i < OFFLINE_BATCH; i++) { |
|
|
|
|
char text[32]; |
|
|
|
|
snprintf(text, sizeof(text), "offline-%02d", i); |
|
|
|
|
if (dm_send(A, g_sh.conv_id, "text", (const uint8_t*)text, (uint32_t)strlen(text)) != 0) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "offline batch send failed index=%d", i); |
|
|
|
|
t->result = 2; |
|
|
|
|
uasync_stop(t->ua); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
t->phase = P_WAIT_MAIL; |
|
|
|
|
break; |
|
|
|
|
case P_WAIT_MAIL: |
|
|
|
|
if (dm_mail_count(C) > 0) { |
|
|
|
|
if (dm_mail_count(C) == OFFLINE_BATCH) { |
|
|
|
|
utun_instance_destroy(A); |
|
|
|
|
t->inst[IDX_A] = NULL; |
|
|
|
|
char cfg[512]; |
|
|
|
|
snprintf(cfg, sizeof(cfg), "%s/a.conf", getenv("UTUN_TEST_DIR")); |
|
|
|
|
t->inst[IDX_A] = utun_instance_create(t->ua, cfg); |
|
|
|
|
if (!t->inst[IDX_A] || utun_instance_init(t->inst[IDX_A]) != 0 || outbox_count(t->inst[IDX_A]) != 1) { |
|
|
|
|
if (!t->inst[IDX_A] || utun_instance_init(t->inst[IDX_A]) != 0 || outbox_count(t->inst[IDX_A]) != OFFLINE_BATCH) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "outbox restart failed"); |
|
|
|
|
t->result = 2; |
|
|
|
|
uasync_stop(t->ua); |
|
|
|
|
@ -489,7 +509,7 @@ static void dm_tick(void* arg) {
|
|
|
|
|
break; |
|
|
|
|
} |
|
|
|
|
case P_WAIT_B2: |
|
|
|
|
if (dm_msg_has(B, "hello2") && dm_mail_count(C) == 0 && outbox_count(A) == 0) t->phase = P_DONE; |
|
|
|
|
if (offline_batch_received(B) && dm_mail_count(C) == 0 && outbox_count(A) == 0) t->phase = P_DONE; |
|
|
|
|
break; |
|
|
|
|
case P_DONE: |
|
|
|
|
t->result = 1; |
|
|
|
|
|