|
|
|
|
@ -591,3 +591,404 @@ void dm_send_file_trampoline(void* arg) {
|
|
|
|
|
dm_send_file(req->inst, req->conv, req->path, req->name); |
|
|
|
|
u_free(req); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* STORE сначала резервирует место и фиксирует expiry; повтор не продлевает TTL. */ |
|
|
|
|
static void dm_custody_store(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* p, size_t len) { |
|
|
|
|
if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: short STORE"); return; } |
|
|
|
|
uint64_t recipient, author, size; |
|
|
|
|
memcpy(&recipient, p, 8); p += 8; len -= 8; |
|
|
|
|
const uint8_t* desc = dm_message_media(p, len); |
|
|
|
|
if (!desc || dm_verify_message(inst, recipient, p, len)) return; |
|
|
|
|
memcpy(&author, p + 24, 8); memcpy(&size, desc + 16, 8); |
|
|
|
|
if (!dm_media_scope(inst, group, source, "supernode", author, recipient) || |
|
|
|
|
!dm_media_scope(inst, group, inst->node_id, "storage", author, recipient)) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized STORE source=%llu group=%llu", |
|
|
|
|
(unsigned long long)source, (unsigned long long)group); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
uint8_t hash[32]; SHA256(p, len, hash); |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body_hash,receipt FROM dm_custody_done WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
if (rc == SQLITE_ROW) { |
|
|
|
|
int exact = sqlite3_column_bytes(st, 0) == 32 && !memcmp(sqlite3_column_blob(st, 0), hash, 32); |
|
|
|
|
uint8_t receipt[DM_MEDIA_RECEIPT_SIZE]; |
|
|
|
|
int delivered = sqlite3_column_bytes(st, 1) == sizeof(receipt); |
|
|
|
|
if (delivered) memcpy(receipt, sqlite3_column_blob(st, 1), sizeof(receipt)); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: STORE conflicts with retired UUID"); return; } |
|
|
|
|
dm_media_send(inst, group, source, delivered ? DM_MEDIA_RECEIPT : DM_MEDIA_EXPIRED, |
|
|
|
|
delivered ? receipt : desc, delivered ? sizeof(receipt) : 16); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: tombstone lookup failed rc=%d", rc); return; } |
|
|
|
|
st = dm_media_prepare(inst, "SELECT body,ready,expires FROM dm_custody WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); |
|
|
|
|
rc = sqlite3_step(st); |
|
|
|
|
if (rc == SQLITE_ROW) { |
|
|
|
|
int exact = sqlite3_column_bytes(st, 0) == (int)len && !memcmp(sqlite3_column_blob(st, 0), p, len); |
|
|
|
|
int ready = sqlite3_column_int(st, 1); |
|
|
|
|
uint64_t expiry = (uint64_t)sqlite3_column_int64(st, 2); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: conflicting custody body"); return; } |
|
|
|
|
if (ready == 1 && expiry > ntp_time_get_seconds(inst)) { |
|
|
|
|
uint8_t stored[24]; memcpy(stored, desc, 16); memcpy(stored + 16, &expiry, 8); |
|
|
|
|
dm_media_send(inst, group, source, DM_MEDIA_STORED, stored, sizeof(stored)); |
|
|
|
|
} |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: custody lookup failed rc=%d", rc); return; } |
|
|
|
|
st = dm_media_prepare(inst, "SELECT COALESCE(sum(cipher_size),0) FROM dm_custody"); |
|
|
|
|
if (!st) return; |
|
|
|
|
rc = sqlite3_step(st); |
|
|
|
|
uint64_t used = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 0) : UINT64_MAX; |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
uint64_t quota = (uint64_t)chat_setting_get_int(inst, "dm_media_storage_mb", 1024) * 1024 * 1024; |
|
|
|
|
uint64_t cipher_size = dm_media_cipher_size(size); |
|
|
|
|
if (used > quota || cipher_size > quota - used) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: custody quota refused used=%llu requested=%llu quota=%llu", |
|
|
|
|
(unsigned long long)used, (unsigned long long)cipher_size, (unsigned long long)quota); |
|
|
|
|
uint8_t reject[17]; memcpy(reject, desc, 16); reject[16] = 1; |
|
|
|
|
dm_media_send(inst, group, source, DM_MEDIA_REJECT, reject, sizeof(reject)); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
uint64_t expires = ntp_time_get_seconds(inst) + (uint64_t)chat_setting_get_int(inst, "dm_media_ttl_sec", 604800); |
|
|
|
|
st = dm_media_prepare(inst, "INSERT INTO dm_custody(id,author,recipient,plain_size,cipher_size,body,body_hash," |
|
|
|
|
"group_id,supernode,expires) VALUES(?,?,?,?,?,?,?,?,?,?)"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, desc, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)author); |
|
|
|
|
sqlite3_bind_int64(st, 3, (sqlite3_int64)recipient); sqlite3_bind_int64(st, 4, (sqlite3_int64)size); |
|
|
|
|
sqlite3_bind_int64(st, 5, (sqlite3_int64)cipher_size); sqlite3_bind_blob(st, 6, p, (int)len, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_blob(st, 7, hash, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 8, (sqlite3_int64)group); |
|
|
|
|
sqlite3_bind_int64(st, 9, (sqlite3_int64)source); sqlite3_bind_int64(st, 10, (sqlite3_int64)expires); |
|
|
|
|
if (!dm_media_step(inst, st)) |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: custody reserved author=%llu recipient=%llu bytes=%llu expires=%llu", |
|
|
|
|
(unsigned long long)author, (unsigned long long)recipient, (unsigned long long)cipher_size, (unsigned long long)expires); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Удаление подтверждается только после durable удаления файла и копии из БД. */ |
|
|
|
|
static void dm_custody_remove(struct dm_media_state* media, const uint8_t id[16]) { |
|
|
|
|
struct UTUN_INSTANCE* inst = media->inst; |
|
|
|
|
if (dm_job_find(media, id, 1)) return; |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body_hash,receipt,supernode,group_id FROM dm_custody WHERE id=? AND ready=-1"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
uint8_t hash[32], receipt[DM_MEDIA_RECEIPT_SIZE]; |
|
|
|
|
if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != sizeof(hash)) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: delete state missing/corrupt rc=%d", rc); |
|
|
|
|
sqlite3_finalize(st); return; |
|
|
|
|
} |
|
|
|
|
memcpy(hash, sqlite3_column_blob(st, 0), sizeof(hash)); |
|
|
|
|
int delivered = sqlite3_column_bytes(st, 1) == sizeof(receipt); |
|
|
|
|
if (delivered) memcpy(receipt, sqlite3_column_blob(st, 1), sizeof(receipt)); |
|
|
|
|
uint64_t super = (uint64_t)sqlite3_column_int64(st, 2), group = (uint64_t)sqlite3_column_int64(st, 3); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
char path[1024]; |
|
|
|
|
const char* suffixes[] = {".cipher", ".download"}; |
|
|
|
|
for (unsigned i = 0; i < sizeof(suffixes) / sizeof(suffixes[0]); i++) |
|
|
|
|
if (dm_file_path(media, 1, id, suffixes[i], path, sizeof(path)) || dm_file_remove(path)) return; |
|
|
|
|
if (dm_dir_sync(path) || dm_media_exec(inst, "BEGIN IMMEDIATE")) return; |
|
|
|
|
st = dm_media_prepare(inst, "INSERT INTO dm_custody_done(id,body_hash,receipt) VALUES(?,?,?)" |
|
|
|
|
" ON CONFLICT(id) DO UPDATE SET receipt=COALESCE(excluded.receipt,receipt)"); |
|
|
|
|
if (!st) goto rollback; |
|
|
|
|
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); sqlite3_bind_blob(st, 2, hash, 32, SQLITE_STATIC); |
|
|
|
|
if (delivered) sqlite3_bind_blob(st, 3, receipt, sizeof(receipt), SQLITE_STATIC); else sqlite3_bind_null(st, 3); |
|
|
|
|
if (dm_media_step(inst, st)) goto rollback; |
|
|
|
|
st = dm_media_prepare(inst, "DELETE FROM dm_custody WHERE id=? AND ready=-1"); |
|
|
|
|
if (!st) goto rollback; |
|
|
|
|
sqlite3_bind_blob(st, 1, id, 16, SQLITE_STATIC); |
|
|
|
|
if (dm_media_step(inst, st) || dm_media_exec(inst, "COMMIT")) goto rollback; |
|
|
|
|
dm_media_send(inst, group, super, delivered ? DM_MEDIA_DELETE_ACK : DM_MEDIA_EXPIRED, id, 16); |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: custody copy removed reason=%s super=%llu", |
|
|
|
|
delivered ? "delivered" : "expired", (unsigned long long)super); |
|
|
|
|
return; |
|
|
|
|
rollback: |
|
|
|
|
dm_media_exec(inst, "ROLLBACK"); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Подпись B позволяет любому авторизованному S завершить удаление точного артефакта. */ |
|
|
|
|
static void dm_custody_delete(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* r) { |
|
|
|
|
if (dm_media_verify_receipt(inst, r, NULL, 0)) return; |
|
|
|
|
uint64_t author, recipient; |
|
|
|
|
memcpy(&author, r + 24, 8); memcpy(&recipient, r + 32, 8); |
|
|
|
|
if (!dm_media_scope(inst, group, source, "supernode", author, recipient) || |
|
|
|
|
!dm_media_scope(inst, group, inst->node_id, "storage", author, recipient)) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized DELETE source=%llu", (unsigned long long)source); return; |
|
|
|
|
} |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT body FROM dm_custody WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
if (rc == SQLITE_ROW) { |
|
|
|
|
int valid = !dm_media_verify_receipt(inst, r, sqlite3_column_blob(st, 0), (size_t)sqlite3_column_bytes(st, 0)); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (!valid) return; |
|
|
|
|
st = dm_media_prepare(inst, "UPDATE dm_custody SET ready=-1,receipt=?,retry_at=0 WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, r, DM_MEDIA_RECEIPT_SIZE, SQLITE_STATIC); sqlite3_bind_blob(st, 2, r + 40, 16, SQLITE_STATIC); |
|
|
|
|
if (dm_media_step(inst, st)) return; |
|
|
|
|
file_transfer_cancel_object(inst->file_transfer, r + 40); |
|
|
|
|
dm_custody_remove(inst->dm_media, r + 40); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: delete lookup failed rc=%d", rc); return; } |
|
|
|
|
st = dm_media_prepare(inst, "SELECT body_hash FROM dm_custody_done WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); |
|
|
|
|
rc = sqlite3_step(st); |
|
|
|
|
int exact = rc == SQLITE_DONE || (rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == 32 && |
|
|
|
|
!memcmp(sqlite3_column_blob(st, 0), r + 56, 32)); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (!exact) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: delete conflicts with retired artifact"); return; } |
|
|
|
|
st = dm_media_prepare(inst, "INSERT INTO dm_custody_done(id,body_hash,receipt) VALUES(?,?,?)" |
|
|
|
|
" ON CONFLICT(id) DO UPDATE SET receipt=excluded.receipt"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, r + 40, 16, SQLITE_STATIC); sqlite3_bind_blob(st, 2, r + 56, 32, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_blob(st, 3, r, DM_MEDIA_RECEIPT_SIZE, SQLITE_STATIC); |
|
|
|
|
if (!dm_media_step(inst, st)) dm_media_send(inst, group, source, DM_MEDIA_DELETE_ACK, r + 40, 16); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* OFFER не меняет идентичность беседы. Чужой holder допустим только от S. */ |
|
|
|
|
static void dm_media_offer(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, const uint8_t* p, size_t len) { |
|
|
|
|
if (len <= 24) { DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: short OFFER"); return; } |
|
|
|
|
uint64_t recipient, holder, expiry, author; |
|
|
|
|
memcpy(&recipient, p, 8); memcpy(&holder, p + 8, 8); memcpy(&expiry, p + 16, 8); |
|
|
|
|
p += 24; len -= 24; |
|
|
|
|
const uint8_t* desc = dm_message_media(p, len); |
|
|
|
|
if (recipient != inst->node_id || !desc || dm_verify_message(inst, recipient, p, len)) return; |
|
|
|
|
memcpy(&author, p + 24, 8); |
|
|
|
|
if ((source == author && (holder != author || expiry)) || |
|
|
|
|
(source != author && !dm_media_scope(inst, group, source, "supernode", author, recipient)) || |
|
|
|
|
(holder != author && !dm_media_scope(inst, group, holder, "storage", author, recipient))) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unauthorized OFFER source=%llu holder=%llu", |
|
|
|
|
(unsigned long long)source, (unsigned long long)holder); return; |
|
|
|
|
} |
|
|
|
|
uint8_t receipt[DM_RECEIPT_SIZE]; |
|
|
|
|
if (dm_accept_message(inst, p, len, receipt)) return; |
|
|
|
|
dm_mailbox_receipt(inst, receipt); |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "UPDATE dm_files SET holder=?,expires=?,group_id=?,retry_at=0 WHERE id=? AND ready<=0"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)holder); sqlite3_bind_int64(st, 2, (sqlite3_int64)expiry); |
|
|
|
|
sqlite3_bind_int64(st, 3, (sqlite3_int64)group); sqlite3_bind_blob(st, 4, desc, 16, SQLITE_STATIC); |
|
|
|
|
dm_media_step(inst, st); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Строгие границы control, независимые подписи DM и квитанций. */ |
|
|
|
|
static void dm_media_control(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t source, |
|
|
|
|
uint8_t cmd, const uint8_t* p, size_t len) { |
|
|
|
|
if (!inst->dm_media || inst->dm_media->closing) return; |
|
|
|
|
if (cmd == DM_MEDIA_OFFER) dm_media_offer(inst, group, source, p, len); |
|
|
|
|
else if (cmd == DM_MEDIA_STORE) dm_custody_store(inst, group, source, p, len); |
|
|
|
|
else if (cmd == DM_MEDIA_DELETE && len == DM_MEDIA_RECEIPT_SIZE) dm_custody_delete(inst, group, source, p); |
|
|
|
|
else if (cmd == DM_MEDIA_RECEIPT && len == DM_MEDIA_RECEIPT_SIZE) { |
|
|
|
|
uint64_t author; memcpy(&author, p + 24, 8); |
|
|
|
|
if (author == inst->node_id) { |
|
|
|
|
if (!dm_media_accept_receipt(inst, p)) dm_mailbox_media_receipt(inst, p); |
|
|
|
|
} else dm_mailbox_media_recv(dm_mailbox_get_media(inst), group, source, cmd, p, len); |
|
|
|
|
} else if ((cmd == DM_MEDIA_STORED && len == 24) || (cmd == DM_MEDIA_DELETE_ACK && len == 16) || |
|
|
|
|
(cmd == DM_MEDIA_EXPIRED && len == 16) || (cmd == DM_MEDIA_REJECT && len == 17)) { |
|
|
|
|
dm_mailbox_media_recv(dm_mailbox_get_media(inst), group, source, cmd, p, len); |
|
|
|
|
if (cmd == DM_MEDIA_EXPIRED) { |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT author,recipient FROM dm_files WHERE id=? AND dir=0 AND ready<=0"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
uint64_t author = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 0) : 0; |
|
|
|
|
uint64_t recipient = rc == SQLITE_ROW ? (uint64_t)sqlite3_column_int64(st, 1) : 0; |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
if (author && dm_media_scope(inst, group, source, "supernode", author, recipient)) { |
|
|
|
|
st = dm_media_prepare(inst, "UPDATE dm_files SET ready=-1,expires=0 WHERE id=? AND ready<=0"); |
|
|
|
|
if (st) { sqlite3_bind_blob(st, 1, p, 16, SQLITE_STATIC); dm_media_step(inst, st); } |
|
|
|
|
dm_media_notify(inst, dm_derive_conv_id(author, recipient)); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} else DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: invalid control cmd=%u bytes=%zu", cmd, len); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void dm_media_recv(struct ETCP_CONN* conn, struct ll_entry* e) { |
|
|
|
|
if (!conn || !e || !e->dgram || e->len <= ROUTER_SVC_PAYLOAD_OFF) goto done; |
|
|
|
|
uint64_t group, source; |
|
|
|
|
memcpy(&group, e->dgram + ROUTER_SVC_GROUP_OFF, 8); memcpy(&source, e->dgram + ROUTER_SVC_SRC_OFF, 8); |
|
|
|
|
if (!(e->dgram[ROUTER_SVC_FLAGS_OFF] & ROUTER_FLAG_SIGNED)) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DM, "dm_media: unsigned control source=%llu", (unsigned long long)source); goto done; |
|
|
|
|
} |
|
|
|
|
const uint8_t* p = e->dgram + ROUTER_SVC_PAYLOAD_OFF; |
|
|
|
|
dm_media_control(conn->instance, group, source, p[0], p + 1, e->len - ROUTER_SVC_PAYLOAD_OFF - 1); |
|
|
|
|
done: |
|
|
|
|
if (e) { queue_dgram_free(e); queue_entry_free(e); } |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Текст и файл имеют независимые durable retries; наличие text ACK не останавливает медиа. */ |
|
|
|
|
static void dm_files_tick(struct dm_media_state* media, uint64_t now) { |
|
|
|
|
struct UTUN_INSTANCE* inst = media->inst; |
|
|
|
|
for (int i = 0; i < 8; i++) { |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT id,dir,author,recipient,group_id,body,ready,holder,expires,receipt" |
|
|
|
|
" FROM dm_files WHERE retry_at<=? ORDER BY retry_at,seq LIMIT 1"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)now); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
if (rc == SQLITE_DONE) { sqlite3_finalize(st); break; } |
|
|
|
|
if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != 16) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: local files scan failed rc=%d", rc); sqlite3_finalize(st); return; |
|
|
|
|
} |
|
|
|
|
uint8_t id[16], body[DM_MEDIA_BODY_MAX], receipt[DM_MEDIA_RECEIPT_SIZE]; |
|
|
|
|
memcpy(id, sqlite3_column_blob(st, 0), 16); |
|
|
|
|
int dir = sqlite3_column_int(st, 1), ready = sqlite3_column_int(st, 6); |
|
|
|
|
uint64_t author = (uint64_t)sqlite3_column_int64(st, 2), recipient = (uint64_t)sqlite3_column_int64(st, 3); |
|
|
|
|
uint64_t group = (uint64_t)sqlite3_column_int64(st, 4), holder = (uint64_t)sqlite3_column_int64(st, 7); |
|
|
|
|
uint64_t expiry = (uint64_t)sqlite3_column_int64(st, 8); |
|
|
|
|
size_t len = (size_t)sqlite3_column_bytes(st, 5); |
|
|
|
|
int delivered = sqlite3_column_bytes(st, 9) == sizeof(receipt); |
|
|
|
|
if (len <= sizeof(body)) memcpy(body, sqlite3_column_blob(st, 5), len); |
|
|
|
|
if (delivered) memcpy(receipt, sqlite3_column_blob(st, 9), sizeof(receipt)); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
st = dm_media_prepare(inst, "UPDATE dm_files SET retry_at=? WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)(now + (dir && delivered ? 3600 : 5))); |
|
|
|
|
sqlite3_bind_blob(st, 2, id, 16, SQLITE_STATIC); |
|
|
|
|
if (dm_media_step(inst, st)) return; |
|
|
|
|
const uint8_t* desc = len <= sizeof(body) ? dm_message_media(body, len) : NULL; |
|
|
|
|
if (!desc) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: corrupt saved file descriptor"); continue; } |
|
|
|
|
if (dir && !delivered) { |
|
|
|
|
uint64_t gid = group; |
|
|
|
|
if (dm_route_group(inst, recipient, group, &gid)) { |
|
|
|
|
uint8_t offer[24 + DM_MEDIA_BODY_MAX]; |
|
|
|
|
memcpy(offer, &recipient, 8); memcpy(offer + 8, &author, 8); memset(offer + 16, 0, 8); |
|
|
|
|
memcpy(offer + 24, body, len); |
|
|
|
|
dm_media_send(inst, gid, recipient, DM_MEDIA_OFFER, offer, len + 24); |
|
|
|
|
} else dm_mailbox_put(inst, recipient, author, body, len); |
|
|
|
|
} else if (!dir && ready == 1 && delivered) { |
|
|
|
|
uint64_t gid = group; |
|
|
|
|
if (dm_route_group(inst, author, group, &gid)) dm_media_send(inst, gid, author, DM_MEDIA_RECEIPT, receipt, sizeof(receipt)); |
|
|
|
|
dm_mailbox_media_receipt(inst, receipt); |
|
|
|
|
} else if (!dir && ready <= 0 && !dm_job_find(media, id, 0)) { |
|
|
|
|
uint64_t peer = author, gid = group; |
|
|
|
|
if (!dm_route_group(inst, author, group, &gid)) { |
|
|
|
|
peer = holder; |
|
|
|
|
if (holder == author || (expiry && expiry <= now) || !dm_route_group(inst, holder, group, &gid)) continue; |
|
|
|
|
} |
|
|
|
|
dm_file_pull(media, 0, gid, peer, desc, author); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Периодическая уборка TTL не вытесняет подтверждённые файлы раньше их expiry. */ |
|
|
|
|
static void dm_custody_tick(struct dm_media_state* media, uint64_t now) { |
|
|
|
|
struct UTUN_INSTANCE* inst = media->inst; |
|
|
|
|
for (int i = 0; i < 8; i++) { |
|
|
|
|
sqlite3_stmt* st = dm_media_prepare(inst, "SELECT id,body,author,supernode,group_id,expires,ready FROM dm_custody" |
|
|
|
|
" WHERE retry_at<=? ORDER BY retry_at LIMIT 1"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)now); |
|
|
|
|
int rc = sqlite3_step(st); |
|
|
|
|
if (rc == SQLITE_DONE) { sqlite3_finalize(st); break; } |
|
|
|
|
if (rc != SQLITE_ROW || sqlite3_column_bytes(st, 0) != 16) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: custody scan failed rc=%d", rc); sqlite3_finalize(st); return; |
|
|
|
|
} |
|
|
|
|
uint8_t id[16], body[DM_MEDIA_BODY_MAX]; |
|
|
|
|
memcpy(id, sqlite3_column_blob(st, 0), 16); |
|
|
|
|
size_t len = (size_t)sqlite3_column_bytes(st, 1); |
|
|
|
|
if (len <= sizeof(body)) memcpy(body, sqlite3_column_blob(st, 1), len); |
|
|
|
|
uint64_t author = (uint64_t)sqlite3_column_int64(st, 2), super = (uint64_t)sqlite3_column_int64(st, 3); |
|
|
|
|
uint64_t group = (uint64_t)sqlite3_column_int64(st, 4), expires = (uint64_t)sqlite3_column_int64(st, 5); |
|
|
|
|
int ready = sqlite3_column_int(st, 6); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
st = dm_media_prepare(inst, "UPDATE dm_custody SET retry_at=?,ready=CASE WHEN expires<=? THEN -1 ELSE ready END WHERE id=?"); |
|
|
|
|
if (!st) return; |
|
|
|
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)(now + 5)); sqlite3_bind_int64(st, 2, (sqlite3_int64)now); |
|
|
|
|
sqlite3_bind_blob(st, 3, id, 16, SQLITE_STATIC); |
|
|
|
|
if (dm_media_step(inst, st)) return; |
|
|
|
|
if (expires <= now || ready < 0) { |
|
|
|
|
file_transfer_cancel_object(inst->file_transfer, id); |
|
|
|
|
dm_custody_remove(media, id); |
|
|
|
|
} else if (!ready) { |
|
|
|
|
const uint8_t* desc = len <= sizeof(body) ? dm_message_media(body, len) : NULL; |
|
|
|
|
uint64_t gid = group; |
|
|
|
|
if (desc && dm_route_group(inst, author, group, &gid)) dm_file_pull(media, 1, gid, author, desc, author); |
|
|
|
|
} else { |
|
|
|
|
uint8_t stored[24]; memcpy(stored, id, 16); memcpy(stored + 16, &expires, 8); |
|
|
|
|
dm_media_send(inst, group, super, DM_MEDIA_STORED, stored, sizeof(stored)); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void dm_media_tick(void* arg) { |
|
|
|
|
struct dm_media_state* media = arg; |
|
|
|
|
media->timer = NULL; |
|
|
|
|
uint64_t now = ntp_time_get_seconds(media->inst); |
|
|
|
|
dm_files_tick(media, now); |
|
|
|
|
dm_custody_tick(media, now); |
|
|
|
|
media->timer = uasync_set_timeout(media->inst->ua, 10000, media, dm_media_tick, "dm_media"); |
|
|
|
|
if (!media->timer) DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: retry timer allocation failed"); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
int dm_media_init(struct UTUN_INSTANCE* inst) { |
|
|
|
|
if (!inst || !inst->dm || !inst->media_async || inst->dm_media) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: invalid initialization state"); return -1; |
|
|
|
|
} |
|
|
|
|
struct dm_media_state* media = u_calloc(1, sizeof(*media)); |
|
|
|
|
if (!media) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: owner allocation failed"); return -1; } |
|
|
|
|
media->inst = inst; |
|
|
|
|
const char* base = inst->config->global.db_path; |
|
|
|
|
if (!base[0]) base = "."; |
|
|
|
|
int n = snprintf(media->root, sizeof(media->root), "%s/dm_media", base); |
|
|
|
|
int m = snprintf(media->pending, sizeof(media->pending), "%s/dm_pending", base); |
|
|
|
|
if (n < 0 || (size_t)n >= sizeof(media->root) || m < 0 || (size_t)m >= sizeof(media->pending) || |
|
|
|
|
(utun_mkdir(media->root, 0700) && errno != EEXIST) || (utun_mkdir(media->pending, 0700) && errno != EEXIST)) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: cannot create file roots: %s", strerror(errno)); goto fail; |
|
|
|
|
} |
|
|
|
|
if (dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_files(id BLOB PRIMARY KEY,conv_id TEXT NOT NULL,dir INTEGER NOT NULL," |
|
|
|
|
"seq INTEGER NOT NULL,author INTEGER NOT NULL,recipient INTEGER NOT NULL,plain_size INTEGER NOT NULL," |
|
|
|
|
"cipher_hash BLOB NOT NULL,body BLOB NOT NULL,group_id INTEGER NOT NULL,local_path TEXT NOT NULL," |
|
|
|
|
"ready INTEGER NOT NULL DEFAULT 0,holder INTEGER NOT NULL,expires INTEGER NOT NULL DEFAULT 0," |
|
|
|
|
"receipt BLOB,downloading INTEGER NOT NULL DEFAULT 0,retry_at INTEGER NOT NULL DEFAULT 0," |
|
|
|
|
"UNIQUE(conv_id,dir,seq))") || |
|
|
|
|
dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_custody(id BLOB PRIMARY KEY,author INTEGER NOT NULL,recipient INTEGER NOT NULL," |
|
|
|
|
"plain_size INTEGER NOT NULL,cipher_size INTEGER NOT NULL,body BLOB NOT NULL,body_hash BLOB NOT NULL," |
|
|
|
|
"group_id INTEGER NOT NULL,supernode INTEGER NOT NULL,expires INTEGER NOT NULL," |
|
|
|
|
"ready INTEGER NOT NULL DEFAULT 0,receipt BLOB,retry_at INTEGER NOT NULL DEFAULT 0)") || |
|
|
|
|
dm_media_exec(inst, "CREATE TABLE IF NOT EXISTS dm_custody_done(id BLOB PRIMARY KEY,body_hash BLOB NOT NULL,receipt BLOB)") || |
|
|
|
|
dm_media_exec(inst, "UPDATE dm_files SET downloading=0,retry_at=0") || dm_media_exec(inst, "UPDATE dm_custody SET retry_at=0")) goto fail; |
|
|
|
|
inst->dm_media = media; |
|
|
|
|
if (!file_transfer_create(inst, dm_file_lookup, media) || etcp_router_bind(inst, ETCP_RT_ID_DM_MEDIA, dm_media_recv)) goto fail; |
|
|
|
|
media->timer = uasync_set_timeout(inst->ua, 10000, media, dm_media_tick, "dm_media"); |
|
|
|
|
if (!media->timer) goto fail; |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: initialized root=%s custody=%s ttl=%ds quota=%dMiB", media->root, media->pending, |
|
|
|
|
chat_setting_get_int(inst, "dm_media_ttl_sec", 604800), chat_setting_get_int(inst, "dm_media_storage_mb", 1024)); |
|
|
|
|
return 0; |
|
|
|
|
fail: |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: initialization failed"); |
|
|
|
|
if (inst->dm_media) dm_media_destroy(inst); else u_free(media); |
|
|
|
|
return -1; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Снять producers до ожидания workers; callbacks отмены ещё имеют живого owner. */ |
|
|
|
|
void dm_media_quiesce(struct UTUN_INSTANCE* inst) { |
|
|
|
|
struct dm_media_state* media = inst ? inst->dm_media : NULL; |
|
|
|
|
if (!media || media->closing) return; |
|
|
|
|
media->closing = 1; |
|
|
|
|
if (media->timer) uasync_cancel_timeout(inst->ua, media->timer); |
|
|
|
|
media->timer = NULL; |
|
|
|
|
etcp_router_unbind(inst, ETCP_RT_ID_DM_MEDIA); |
|
|
|
|
file_transfer_destroy(inst->file_transfer); |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: quiesced, waiting for workers"); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void dm_media_destroy(struct UTUN_INSTANCE* inst) { |
|
|
|
|
if (!inst || !inst->dm_media) return; |
|
|
|
|
dm_media_quiesce(inst); |
|
|
|
|
struct dm_media_state* media = inst->dm_media; |
|
|
|
|
if (media->jobs) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_DM, "dm_media: destroy before media_async workers drained"); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
inst->dm_media = NULL; |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DM, "dm_media: destroyed, durable files/jobs retained"); |
|
|
|
|
u_free(media); |
|
|
|
|
} |
|
|
|
|
|