From 11e98b300e2c91d93253c4018baf3b5100245d4e Mon Sep 17 00:00:00 2001 From: ayrisdev Date: Fri, 4 Sep 2026 19:57:00 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20analiz=20zincirini=20atlay=C4=B1p=20ham?= =?UTF-8?q?=20segmenti=20do=C4=9Frudan=20Telegram'a=20g=C3=B6nderme=20modu?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Yeni "Ham Kayıt Modu" ayarı (send_raw_segments_to_telegram, varsayılan kapalı): açıksa her ham segment tamamlanır tamamlanmaz sinyal analizi/ STT/9:16 render zincirine hiç girmeden doğrudan Telegram'a gönderiliyor ve sunucudan siliniyor. Kullanıcı transkript/9:16 istemediğinde ("sen bana direk videoyu gönder") tüm analiz pipeline'ını atlayıp sadece kayıt+gönder+sil yapan basit bir mod. - RawSegmentStatus'a SENT_TO_TELEGRAM eklendi, RawSegment.filePath nullable yapıldı (aynı ShortVideo.filePath deseni — gönderim sonrası null, dosya diskten silinmiş demek). - api-daemon/src/telegram.ts'e sendTelegramVideo eklendi (native fetch+ FormData+Blob ile multipart upload, ~50MB bot-limiti önceden kontrol ediliyor) — worker/telegram.py'deki aynı fonksiyonun Node karşılığı. - streamIngest.ts: segment tamamlanınca, mod açıksa signal-detection kuyruğuna hiç eklemeden doğrudan gönderiyor (kuyruğa eklemek, aşağıda dosyayı silmenin sinyal analizini bozmasına sebep olurdu). - Panel: segment listelerinde filePath=null durumunda video/indir yerine "Telegram'a gönderildi" notu (segments, channels/[id]). Co-Authored-By: Claude Sonnet 5 --- apps/api-daemon/src/capture/streamIngest.ts | 44 ++++++++++++--- apps/api-daemon/src/cleanup.ts | 2 +- apps/api-daemon/src/telegram.ts | 43 ++++++++++++++ apps/frontend/app/actions.ts | 14 ++++- .../app/api/segments/[id]/video/route.ts | 2 +- apps/frontend/app/channels/[id]/page.tsx | 56 ++++++++++++------- apps/frontend/app/segments/page.tsx | 42 ++++++++------ apps/frontend/app/settings/page.tsx | 31 +++++++++- .../migration.sql | 5 ++ packages/db/prisma/schema.prisma | 6 +- 10 files changed, 193 insertions(+), 52 deletions(-) create mode 100644 packages/db/prisma/migrations/20260904000000_raw_segment_telegram_send/migration.sql diff --git a/apps/api-daemon/src/capture/streamIngest.ts b/apps/api-daemon/src/capture/streamIngest.ts index 80a9dd1..3185542 100644 --- a/apps/api-daemon/src/capture/streamIngest.ts +++ b/apps/api-daemon/src/capture/streamIngest.ts @@ -1,5 +1,5 @@ import { spawn, type ChildProcess } from "node:child_process"; -import { mkdir, readFile } from "node:fs/promises"; +import { mkdir, readFile, unlink } from "node:fs/promises"; import path from "node:path"; import { Worker, type Job } from "bullmq"; import { prisma } from "@streamclipper/db"; @@ -7,6 +7,7 @@ import { env } from "../env"; import { QUEUE_NAMES, signalDetectionQueue, type StreamIngestJob } from "../queues"; import { redisConnection } from "../redis"; import { ytdlpAntiBotArgs } from "../ytdlpCookies"; +import { sendTelegramVideo } from "../telegram"; const SEGMENT_LIST_POLL_MS = 5_000; @@ -59,13 +60,20 @@ function parseSegmentListLines(raw: string): SegmentListEntry[] { }); } +async function sendRawSegmentsToTelegramEnabled(): Promise { + const setting = await prisma.appSetting.findUnique({ where: { key: "send_raw_segments_to_telegram" } }); + return setting?.value === "true"; +} + async function processCompletedSegments( entries: SegmentListEntry[], fromIndex: number, outDir: string, sessionId: string, + channelName: string, ): Promise { let processed = fromIndex; + const sendRawToTelegram = await sendRawSegmentsToTelegramEnabled(); for (let i = fromIndex; i < entries.length; i++) { const entry = entries[i]; @@ -82,13 +90,29 @@ async function processCompletedSegments( }, }); - await signalDetectionQueue.add("analyze-segment", { - rawSegmentId: segment.id, - filePath, - sessionId, - }); + if (sendRawToTelegram) { + // Kullanıcı transkript/9:16 zincirini hiç istemiyor — ham segmenti + // analiz kuyruğuna hiç sokmadan doğrudan gönderiyoruz. Kuyruğa + // eklemiyoruz çünkü aşağıda dosyayı silmek, o segmenti daha sonra + // okumaya çalışacak sinyal analizini bozardı. + const caption = `🎬 *${channelName}* — ${duration}sn ham kayıt`; + const sent = await sendTelegramVideo(filePath, caption); + if (sent) { + await unlink(filePath).catch(() => {}); + await prisma.rawSegment.update({ where: { id: segment.id }, data: { status: "SENT_TO_TELEGRAM", filePath: null } }); + console.log(`[stream-ingest] segment ${entry.fileName} sent to Telegram and removed from disk`); + } else { + console.warn(`[stream-ingest] segment ${entry.fileName} failed to send to Telegram, kept on disk`); + } + } else { + await signalDetectionQueue.add("analyze-segment", { + rawSegmentId: segment.id, + filePath, + sessionId, + }); + console.log(`[stream-ingest] segment ready: ${entry.fileName} -> raw_segment ${segment.id}`); + } - console.log(`[stream-ingest] segment ready: ${entry.fileName} -> raw_segment ${segment.id}`); processed = i + 1; } @@ -97,6 +121,8 @@ async function processCompletedSegments( async function runCapture(job: Job): Promise { const { sessionId, youtubeUrl, segmentTimeSec, channelDbId } = job.data; + const channel = await prisma.channel.findUnique({ where: { id: channelDbId } }); + const channelName = channel?.name ?? "Bilinmeyen kanal"; const outDir = path.join(env.sharedMediaRoot, "raw", sessionId); await mkdir(outDir, { recursive: true }); @@ -150,7 +176,7 @@ async function runCapture(job: Job): Promise { if (!raw) return; const entries = parseSegmentListLines(raw); if (entries.length > lastProcessedIndex) { - lastProcessedIndex = await processCompletedSegments(entries, lastProcessedIndex, outDir, sessionId); + lastProcessedIndex = await processCompletedSegments(entries, lastProcessedIndex, outDir, sessionId, channelName); await prisma.streamSession.update({ where: { id: sessionId }, data: { totalSegments: lastProcessedIndex }, @@ -178,7 +204,7 @@ async function runCapture(job: Job): Promise { const raw = await readFile(segmentListPath, "utf8").catch(() => ""); const entries = parseSegmentListLines(raw); if (entries.length > lastProcessedIndex) { - lastProcessedIndex = await processCompletedSegments(entries, lastProcessedIndex, outDir, sessionId); + lastProcessedIndex = await processCompletedSegments(entries, lastProcessedIndex, outDir, sessionId, channelName); } await prisma.streamSession.update({ diff --git a/apps/api-daemon/src/cleanup.ts b/apps/api-daemon/src/cleanup.ts index 3ca0454..7e14c2a 100644 --- a/apps/api-daemon/src/cleanup.ts +++ b/apps/api-daemon/src/cleanup.ts @@ -32,7 +32,7 @@ async function cleanupOnce(): Promise { for (const segment of staleSegments) { await prisma.candidateSegment.deleteMany({ where: { rawSegmentId: segment.id } }); await prisma.rawSegment.delete({ where: { id: segment.id } }); - await unlink(segment.filePath).catch(() => {}); + if (segment.filePath) await unlink(segment.filePath).catch(() => {}); } if (staleSegments.length > 0) { diff --git a/apps/api-daemon/src/telegram.ts b/apps/api-daemon/src/telegram.ts index 9097d70..d13578d 100644 --- a/apps/api-daemon/src/telegram.ts +++ b/apps/api-daemon/src/telegram.ts @@ -1,6 +1,12 @@ +import { readFile, stat } from "node:fs/promises"; +import path from "node:path"; + const BOT_TOKEN = process.env.TELEGRAM_BOT_TOKEN ?? ""; const CHAT_ID = process.env.TELEGRAM_CHAT_ID ?? ""; +// Telegram's bot-upload limit (no local Bot API server). +const MAX_VIDEO_BYTES = 50 * 1024 * 1024; + export async function sendTelegramMessage(text: string): Promise { if (!BOT_TOKEN || !CHAT_ID) { console.warn("[telegram] TELEGRAM_BOT_TOKEN/TELEGRAM_CHAT_ID not set, skipping notification:", text); @@ -18,3 +24,40 @@ export async function sendTelegramMessage(text: string): Promise { console.error("[telegram] failed to send message:", res.status, await res.text()); } } + +/** Uploads a video file to Telegram. Returns false (never throws) on any failure — the caller decides what that means for the local file. */ +export async function sendTelegramVideo(filePath: string, caption: string): Promise { + if (!BOT_TOKEN || !CHAT_ID) { + console.warn("[telegram] TELEGRAM_BOT_TOKEN/TELEGRAM_CHAT_ID not set, skipping video:", caption); + return false; + } + + const { size } = await stat(filePath); + if (size > MAX_VIDEO_BYTES) { + console.warn(`[telegram] video ${filePath} is ${size} bytes, over Telegram's bot-upload limit — skipping`); + return false; + } + + try { + const buffer = await readFile(filePath); + const form = new FormData(); + form.set("chat_id", CHAT_ID); + form.set("caption", caption); + form.set("parse_mode", "Markdown"); + form.set("video", new Blob([buffer], { type: "video/mp4" }), path.basename(filePath)); + + const res = await fetch(`https://api.telegram.org/bot${BOT_TOKEN}/sendVideo`, { + method: "POST", + body: form, + }); + + if (!res.ok) { + console.error("[telegram] failed to send video:", res.status, await res.text()); + return false; + } + return true; + } catch (err) { + console.error("[telegram] failed to upload video:", err); + return false; + } +} diff --git a/apps/frontend/app/actions.ts b/apps/frontend/app/actions.ts index 9f1426f..58722fd 100644 --- a/apps/frontend/app/actions.ts +++ b/apps/frontend/app/actions.ts @@ -40,7 +40,7 @@ export async function deleteRawSegment(segmentId: string) { await deleteShortsForCandidates(segment.candidates.map((c) => c.id)); await prisma.candidateSegment.deleteMany({ where: { rawSegmentId: segmentId } }); await prisma.rawSegment.delete({ where: { id: segmentId } }); - await unlink(segment.filePath).catch(() => {}); + if (segment.filePath) await unlink(segment.filePath).catch(() => {}); revalidatePath("/segments"); revalidatePath("/", "layout"); @@ -61,7 +61,7 @@ export async function deleteChannel(channelId: string) { }); const rawSegmentIds = sessions.flatMap((s) => s.segments.map((seg) => seg.id)); const candidateIds = sessions.flatMap((s) => s.segments.flatMap((seg) => seg.candidates.map((c) => c.id))); - const filePaths = sessions.flatMap((s) => s.segments.map((seg) => seg.filePath)); + const filePaths = sessions.flatMap((s) => s.segments.map((seg) => seg.filePath)).filter((p): p is string => Boolean(p)); await deleteShortsForCandidates(candidateIds); await prisma.candidateSegment.deleteMany({ where: { rawSegmentId: { in: rawSegmentIds } } }); @@ -175,6 +175,16 @@ export async function updateDeleteAfterTelegramSend(formData: FormData) { revalidatePath("/settings"); } +export async function updateSendRawSegmentsToTelegram(formData: FormData) { + const enabled = formData.get("sendRawSegmentsToTelegram") === "true"; + await prisma.appSetting.upsert({ + where: { key: "send_raw_segments_to_telegram" }, + create: { key: "send_raw_segments_to_telegram", value: String(enabled) }, + update: { value: String(enabled) }, + }); + revalidatePath("/settings"); +} + export async function updateNotificationPrefs(formData: FormData) { const streamStart = formData.get("notifyStreamStart") === "true"; const renderDone = formData.get("notifyRenderDone") === "true"; diff --git a/apps/frontend/app/api/segments/[id]/video/route.ts b/apps/frontend/app/api/segments/[id]/video/route.ts index eecf4a5..865bfbc 100644 --- a/apps/frontend/app/api/segments/[id]/video/route.ts +++ b/apps/frontend/app/api/segments/[id]/video/route.ts @@ -6,7 +6,7 @@ export async function GET(req: NextRequest, { params }: { params: Promise<{ id: const { id } = await params; const segment = await prisma.rawSegment.findUnique({ where: { id } }); - if (!segment) return new Response("Not found", { status: 404 }); + if (!segment?.filePath) return new Response("Not found", { status: 404 }); return streamVideoFile(req, segment.filePath, `${id}.mp4`); } diff --git a/apps/frontend/app/channels/[id]/page.tsx b/apps/frontend/app/channels/[id]/page.tsx index b82b7cf..5aa0875 100644 --- a/apps/frontend/app/channels/[id]/page.tsx +++ b/apps/frontend/app/channels/[id]/page.tsx @@ -30,6 +30,7 @@ const STATUS_BADGE: Record["variant"] PENDING: "muted", PROCESSED: "ok", DISCARDED: "err", + SENT_TO_TELEGRAM: "ok", PENDING_STT: "warn", TRANSCRIBED: "ok", RENDERING: "warn", @@ -225,30 +226,43 @@ export default async function ChannelDetailPage({ params }: { params: Promise<{ {segment.status} - {segment.filePath} - - - -
- -
- - -
- - {segment.candidates.length === 0 && segment.status !== "PENDING" && ( -

- Bu segmentte aday klip bulunamadı (sinyal analizi ilginç bir an tespit etmedi). + {segment.filePath ? ( + <> + {segment.filePath} + +

+ +
+ + +
+ + ) : ( +

+ ✅ Telegram'a gönderildi, sunucuda saklanmıyor.

)} + {segment.status === "SENT_TO_TELEGRAM" ? ( +

+ Analiz edilmeden (transkript/9:16 atlanarak) doğrudan gönderildi. +

+ ) : ( + segment.candidates.length === 0 && + segment.status !== "PENDING" && ( +

+ Bu segmentte aday klip bulunamadı (sinyal analizi ilginç bir an tespit etmedi). +

+ ) + )} + {segment.candidates.length > 0 && (
{segment.candidates.map((c) => ( diff --git a/apps/frontend/app/segments/page.tsx b/apps/frontend/app/segments/page.tsx index 85e8090..69ca7c7 100644 --- a/apps/frontend/app/segments/page.tsx +++ b/apps/frontend/app/segments/page.tsx @@ -13,6 +13,7 @@ const STATUS_BADGE: Record["variant"] PENDING: "muted", PROCESSED: "ok", DISCARDED: "err", + SENT_TO_TELEGRAM: "ok", PENDING_STT: "warn", TRANSCRIBED: "ok", RENDERING: "warn", @@ -53,23 +54,32 @@ export default async function SegmentsPage() { {segment.status} - + {segment.filePath ? ( + <> + +
+ +
+ + +
+ + ) : ( +

✅ Telegram'a gönderildi, sunucuda saklanmıyor.

+ )} -
- -
- - -
- - {segment.candidates.length === 0 ? ( + {segment.status === "SENT_TO_TELEGRAM" ? ( +

+ Bu segment analiz edilmeden (transkript/9:16 atlanarak) doğrudan Telegram'a gönderildi. +

+ ) : segment.candidates.length === 0 ? (

Bu segmentte aday klip bulunamadı.

) : ( segment.candidates.map((c) => ( diff --git a/apps/frontend/app/settings/page.tsx b/apps/frontend/app/settings/page.tsx index 8b52be9..587aa05 100644 --- a/apps/frontend/app/settings/page.tsx +++ b/apps/frontend/app/settings/page.tsx @@ -6,6 +6,7 @@ import { updateSttEnabled, updateAutoRenderEnabled, updateDeleteAfterTelegramSend, + updateSendRawSegmentsToTelegram, addCookieProfile, updateCookieProfile, deleteCookieProfile, @@ -36,6 +37,7 @@ export default async function SettingsPage() { sttSetting, autoRenderSetting, deleteAfterSendSetting, + sendRawSetting, cookieProfiles, ] = await Promise.all([ prisma.appSetting.findUnique({ where: { key: "ytdlp_cookies" } }), @@ -47,6 +49,7 @@ export default async function SettingsPage() { prisma.appSetting.findUnique({ where: { key: "stt_enabled" } }), prisma.appSetting.findUnique({ where: { key: "auto_render_enabled" } }), prisma.appSetting.findUnique({ where: { key: "delete_after_telegram_send" } }), + prisma.appSetting.findUnique({ where: { key: "send_raw_segments_to_telegram" } }), prisma.cookieProfile.findMany({ orderBy: { createdAt: "asc" } }), ]); @@ -60,6 +63,7 @@ export default async function SettingsPage() { const sttEnabled = sttSetting?.value !== "false"; const autoRenderEnabled = autoRenderSetting?.value !== "false"; const deleteAfterTelegramSend = deleteAfterSendSetting?.value === "true"; + const sendRawSegmentsToTelegram = sendRawSetting?.value === "true"; return (
@@ -259,8 +263,33 @@ export default async function SettingsPage() { - Depolama + Ham Kayıt Modu + Açarsan her ham segment (analiz/transkript/9:16 zincirine hiç girmeden) tamamlanır tamamlanmaz doğrudan + Telegram'a gönderilip sunucudan silinir — transkript ve 9:16 render tamamen atlanır. Bunu açtıysan + STT ve 9:16 Render ayarlarının kapalı olması mantıklı (aksi halde ikisi de boşa kalır, çünkü bu mod + segmenti analiz kuyruğuna hiç sokmuyor). ~50MB'ı geçen segmentler gönderilemez, sunucuda kalır — + gerekirse yukarıdaki "Sistem Parametreleri"nden segment süresini kısalt. + + +
+ + + + + + +
+
+ + + + Depolama (9:16 klipler) + + Yukarıdaki "Ham Kayıt Modu"ndan farklı — bu, STT+9:16 render zincirinden geçmiş klipler için. Açarsan her render biten video Telegram'a gönderilip sunucudan kalıcı olarak silinir — sadece Telegram'daki kopyada kalır (geri alınamaz). Gönderim başarısız olursa dosya sunucuda kalır. Kanal sayısı arttıkça disk doluluğunu önlemek için. diff --git a/packages/db/prisma/migrations/20260904000000_raw_segment_telegram_send/migration.sql b/packages/db/prisma/migrations/20260904000000_raw_segment_telegram_send/migration.sql new file mode 100644 index 0000000..863d973 --- /dev/null +++ b/packages/db/prisma/migrations/20260904000000_raw_segment_telegram_send/migration.sql @@ -0,0 +1,5 @@ +-- AlterEnum +ALTER TYPE "RawSegmentStatus" ADD VALUE 'SENT_TO_TELEGRAM'; + +-- AlterTable +ALTER TABLE "raw_segments" ALTER COLUMN "file_path" DROP NOT NULL; diff --git a/packages/db/prisma/schema.prisma b/packages/db/prisma/schema.prisma index 32ec9d0..20b91c1 100644 --- a/packages/db/prisma/schema.prisma +++ b/packages/db/prisma/schema.prisma @@ -11,6 +11,10 @@ enum RawSegmentStatus { PENDING PROCESSED DISCARDED + // Analiz zincirine (sinyal analizi/STT/9:16) hiç girmeden ham haliyle + // doğrudan Telegram'a gönderildi — file_path bu durumda null (dosya + // gönderim sonrası diskten silindi). + SENT_TO_TELEGRAM } enum CandidateSegmentStatus { @@ -60,7 +64,7 @@ model RawSegment { id String @id @default(cuid()) sessionId String @map("session_id") session StreamSession @relation(fields: [sessionId], references: [id]) - filePath String @map("file_path") + filePath String? @map("file_path") duration Int startedAt DateTime @map("started_at") status RawSegmentStatus @default(PENDING)