diff --git a/src/media_async/attachment_send.c b/src/media_async/attachment_send.c new file mode 100644 index 00000000..df381d4d --- /dev/null +++ b/src/media_async/attachment_send.c @@ -0,0 +1,137 @@ +/* Подготовка не обращается к БД; публикация и выбор транспорта выполняются в uasync. */ +#include +#include +#include +#include +#include "../../lib/mem.h" +#include "../../lib/debug_config.h" +#include "../../lib/platform_compat.h" +#include "../../lib/audio_compressor.h" +#include "../utun_instance.h" +#include "../chat/chat_core.h" +#include "../dm/dm_media.h" +#include "../video/video.h" +#include "media_async.h" +#include "voice_file.h" +#include "attachment_send.h" + +/* PCM освобождается в worker после кодирования, до callback. */ +static void attachment_release_pcm(struct attachment_send_req* req) { + if (req->pcm_release) req->pcm_release(req->pcm_owner); + else u_free((void*)req->pcm); + req->pcm = NULL; req->pcm_owner = NULL; req->pcm_release = NULL; + audio_compressor_destroy(req->compressor); req->compressor = NULL; +} + +/* Только пути собственного output; исходный пользовательский файл не удаляется. */ +static void attachment_prepare_work(void* arg) { + struct attachment_send_req* req = arg; + req->error = -1; + char parent[1024]; snprintf(parent, sizeof(parent), "%s", req->output); + for (char* p = parent + 1; *p; p++) { + if (*p != '/') continue; + *p = 0; + int rc = utun_mkdir(parent, 0700); + *p = '/'; + if (rc && errno != EEXIST) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: mkdir failed path=%s", parent); goto done; } + } + if (req->info.kind == ATTACHMENT_VOICE) { + const int16_t* pcm = req->pcm; size_t count = req->pcm_count; + if (req->compressor) { + audio_compressor_flush(req->compressor); + pcm = audio_compressor_output(req->compressor); count = audio_compressor_output_size(req->compressor); + } + req->error = voice_file_encode(pcm, count, 48000, 1, req->preset, req->output, &req->info); + } else if (req->transcode) req->error = 0; /* FFmpeg worker стартует после создания каталога. */ + else { + req->error = ma_copy_file(req->source, req->output); + if (!req->error && req->temporary_source && remove(req->source)) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: source cleanup failed path=%s", req->source); req->error = -1; + } + } +done: + attachment_release_pcm(req); + if (req->error && remove(req->output) && errno != ENOENT) + DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "attachment: partial output cleanup failed path=%s", req->output); +} + +/* Один выбор транспорта после общей подготовки; групповой индекс и PM crypto не смешиваются. */ +static void attachment_publish(struct attachment_send_req* req, int error) { + if (error || req->error || attachment_validate(&req->info)) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: preparation failed target=%s error=%d worker=%d", req->target, error, req->error); + if (remove(req->output) && errno != ENOENT) DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "attachment: output cleanup failed path=%s", req->output); + u_free(req); return; + } + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "attachment: prepared target=%s pm=%d kind=%d duration=%u dimensions=%ux%u", + req->target, req->is_dm, req->info.kind, req->info.duration_ms, req->info.width, req->info.height); + if (req->is_dm) { + struct dm_file_req* file = u_calloc(1, sizeof(*file)); + if (file) { + file->inst = req->inst; file->info = req->info; file->remove_source = 1; + snprintf(file->conv, sizeof(file->conv), "%s", req->target); + snprintf(file->path, sizeof(file->path), "%s", req->output); + dm_send_file_trampoline(file); + } else DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: PM dispatch allocation failed"); + } else { + char text[384]; + if (req->info.kind == ATTACHMENT_VOICE) snprintf(text, sizeof(text), "wf=%s;dur=%u;", req->info.waveform, req->info.duration_ms); + else snprintf(text, sizeof(text), "%s", req->info.name); + size_t length = strlen(text); + struct chat_msg_submit* message = u_calloc(1, sizeof(*message) + length); + if (message) { + message->inst = req->inst; + snprintf(message->channel_id, sizeof(message->channel_id), "%s", req->target); + snprintf(message->content_type, sizeof(message->content_type), "%s", attachment_content_type(req->info.kind)); + snprintf(message->media_src, sizeof(message->media_src), "%s", req->output); + snprintf(message->media_dest, sizeof(message->media_dest), "%s", req->output); + message->duration_ms = req->info.duration_ms; message->width = req->info.width; message->height = req->info.height; + message->data = (uint8_t*)(message + 1); message->data_len = (uint32_t)length; memcpy(message->data, text, length); + chat_core_submit_trampoline(message); + } else DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: channel dispatch allocation failed"); + } + u_free(req); +} + +static void attachment_video_done(void* arg, int error, const struct video_info* info) { + struct attachment_send_req* req = arg; + if (!error && info && info->is_video && info->duration_ms > 0 && info->duration_ms <= 86400000 && + info->width > 0 && info->width <= 16384 && info->height > 0 && info->height <= 16384) { + req->info.duration_ms = (uint32_t)info->duration_ms; + req->info.width = (uint16_t)info->width; req->info.height = (uint16_t)info->height; + } else req->error = -1; + attachment_publish(req, error); +} + +static void attachment_prepare_done(void* arg, int error) { + struct attachment_send_req* req = arg; + if (!error && !req->error && req->transcode) { + if (!video_transcode_start(req->inst->media_async, req->inst->ua, req->source, req->output, attachment_video_done, req)) return; + error = -1; + } + attachment_publish(req, error); +} + +/* Request содержит адресата, выбранного до записи/выбора файла. */ +void attachment_send_trampoline(void* arg) { + struct attachment_send_req* req = arg; + if (!req) return; + if (!req->inst || !req->inst->chat || !req->inst->media_async || !req->target[0] || + strspn(req->target, "0123456789") != strlen(req->target)) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: invalid target/stopped service"); + attachment_release_pcm(req); u_free(req); return; + } + uint8_t uuid[16]; char hex[33]; + if (RAND_bytes(uuid, sizeof(uuid)) != 1) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: UUID generation failed"); attachment_release_pcm(req); u_free(req); return; + } + for (unsigned i = 0; i < 16; i++) snprintf(hex + 2 * i, 3, "%02x", uuid[i]); + const char* base = req->inst->config->global.db_path; + int n = snprintf(req->output, sizeof(req->output), "%s/media/%s/%s.%s", base[0] ? base : ".", + req->is_dm ? "pm-prepared" : req->target, hex, req->info.kind == ATTACHMENT_VOICE ? "opus" : "mp4"); + if (n < 0 || (size_t)n >= sizeof(req->output)) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "attachment: output path too long"); attachment_release_pcm(req); u_free(req); return; + } + if (!req->info.name[0]) snprintf(req->info.name, sizeof(req->info.name), "%s.%s", hex, + req->info.kind == ATTACHMENT_VOICE ? "opus" : "mp4"); + media_async_submit(req->inst->media_async, req->inst->ua, attachment_prepare_work, req, attachment_prepare_done, req); +} diff --git a/src/media_async/attachment_send.h b/src/media_async/attachment_send.h new file mode 100644 index 00000000..ad596d47 --- /dev/null +++ b/src/media_async/attachment_send.h @@ -0,0 +1,31 @@ +/* Подготовить голос/видео и отправить в зафиксированную беседу. Вызывать trampoline + * через uasync_post. Request и PCM/compressor переходят во владение задачи; + * pcm_release освобождает pcm_owner, NULL означает PCM из u_malloc. + * Qt/Android используют общий результат подготовки, доставка остаётся в chat/dm. */ +#ifndef ATTACHMENT_SEND_H +#define ATTACHMENT_SEND_H +#include "attachment.h" +struct UTUN_INSTANCE; +struct audio_compressor; +#ifdef __cplusplus +extern "C" { +#endif +struct attachment_send_req { + struct UTUN_INSTANCE* inst; + char target[64]; + int is_dm; + struct attachment_info info; + char source[1024], output[1024]; + int transcode, temporary_source; + const int16_t* pcm; + size_t pcm_count; + void* pcm_owner; + void (*pcm_release)(void* owner); + struct audio_compressor* compressor; + int preset, error; +}; +void attachment_send_trampoline(void* arg); +#ifdef __cplusplus +} +#endif +#endif