diff --git a/core/queue.py b/core/queue.py
index 07efd82..6c7496f 100644
--- a/core/queue.py
+++ b/core/queue.py
@@ -1,6 +1,7 @@
import os
+import json
import logging
-from typing import Optional
+from typing import Optional, Dict, Any
import redis.asyncio as redis
logger = logging.getLogger(__name__)
@@ -8,9 +9,8 @@ logger = logging.getLogger(__name__)
REDIS_URL = os.getenv("REDIS_URL", "redis://copykar_redis:6379/0" if os.path.exists("/app") else "redis://localhost:6379/0")
class RedisQueue:
- def __init__(self, redis_url: Optional[str] = None, queue_key: str = "copykar:queue:incoming"):
+ def __init__(self, redis_url: Optional[str] = None):
self.redis_url = redis_url or REDIS_URL
- self.queue_key = queue_key
self.client: Optional[redis.Redis] = None
async def connect(self):
@@ -19,22 +19,43 @@ class RedisQueue:
await self.client.ping()
logger.info(f"Connected to Redis at {self.redis_url}")
- async def push(self, post_id: int):
- if not self.client:
- await self.connect()
- await self.client.rpush(self.queue_key, str(post_id))
- logger.info(f"Enqueued post ID {post_id} to Redis queue [{self.queue_key}]")
+ def _get_target_key(self, target_id: int) -> str:
+ return f"copykar:queue:target:{target_id}"
- async def pop(self) -> Optional[int]:
+ async def push_target_post(self, target_id: int, payload: Dict[str, Any]):
if not self.client:
await self.connect()
- val = await self.client.lpop(self.queue_key)
- return int(val) if val else None
+ key = self._get_target_key(target_id)
+ raw_json = json.dumps(payload)
+ await self.client.rpush(key, raw_json)
+ logger.info(f"Enqueued post {payload.get('post_id')} to Target #{target_id} queue [{key}]")
- async def qsize(self) -> int:
+ async def pop_target_post(self, target_id: int) -> Optional[Dict[str, Any]]:
if not self.client:
await self.connect()
- return await self.client.llen(self.queue_key)
+ key = self._get_target_key(target_id)
+ raw = await self.client.lpop(key)
+ if raw:
+ try:
+ return json.loads(raw)
+ except Exception as e:
+ logger.error(f"Error parsing queue JSON from {key}: {e}")
+ return None
+
+ async def get_target_queue_size(self, target_id: int) -> int:
+ if not self.client:
+ await self.connect()
+ key = self._get_target_key(target_id)
+ return await self.client.llen(key)
+
+ async def get_total_queued_posts(self) -> int:
+ if not self.client:
+ await self.connect()
+ keys = await self.client.keys("copykar:queue:target:*")
+ total = 0
+ for k in keys:
+ total += await self.client.llen(k)
+ return total
async def close(self):
if self.client:
diff --git a/db/database.py b/db/database.py
index 77d18e3..d63bb6a 100644
--- a/db/database.py
+++ b/db/database.py
@@ -26,6 +26,9 @@ CREATE TABLE IF NOT EXISTS targets (
post_interval_min INT DEFAULT 30,
personality TEXT DEFAULT '',
custom_footer TEXT DEFAULT '',
+ sleep_start_hour INT DEFAULT 0,
+ sleep_end_hour INT DEFAULT 0,
+ is_sleep_enabled BOOLEAN DEFAULT FALSE,
last_post_time TIMESTAMPTZ,
is_active BOOLEAN DEFAULT TRUE,
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
@@ -71,6 +74,9 @@ CREATE TABLE IF NOT EXISTS settings (
-- Migration safety for existing tables
ALTER TABLE targets ADD COLUMN IF NOT EXISTS personality TEXT DEFAULT '';
ALTER TABLE targets ADD COLUMN IF NOT EXISTS custom_footer TEXT DEFAULT '';
+ALTER TABLE targets ADD COLUMN IF NOT EXISTS sleep_start_hour INT DEFAULT 0;
+ALTER TABLE targets ADD COLUMN IF NOT EXISTS sleep_end_hour INT DEFAULT 0;
+ALTER TABLE targets ADD COLUMN IF NOT EXISTS is_sleep_enabled BOOLEAN DEFAULT FALSE;
ALTER TABLE posts ADD COLUMN IF NOT EXISTS published_to JSONB DEFAULT '[]'::jsonb;
ALTER TABLE posts ADD COLUMN IF NOT EXISTS is_deleted BOOLEAN DEFAULT FALSE;
"""
diff --git a/db/models.py b/db/models.py
index 77200f7..0effd3a 100644
--- a/db/models.py
+++ b/db/models.py
@@ -19,6 +19,9 @@ class TargetChannel:
post_interval_min: int = 30
personality: str = ""
custom_footer: str = ""
+ sleep_start_hour: int = 0
+ sleep_end_hour: int = 0
+ is_sleep_enabled: bool = False
last_post_time: Optional[str] = None
is_active: bool = True
created_at: Optional[str] = None
diff --git a/db/repository.py b/db/repository.py
index e276719..6870ffd 100644
--- a/db/repository.py
+++ b/db/repository.py
@@ -103,6 +103,33 @@ class Repository:
custom_footer, target_id
)
+ async def update_target_schedule(
+ self,
+ target_id: int,
+ post_interval_min: Optional[int] = None,
+ sleep_start_hour: Optional[int] = None,
+ sleep_end_hour: Optional[int] = None,
+ is_sleep_enabled: Optional[bool] = None
+ ) -> None:
+ pool = await self._get_pool()
+ async with pool.acquire() as conn:
+ target = await self.get_target_by_id(target_id)
+ if not target:
+ return
+ new_interval = post_interval_min if post_interval_min is not None else target.post_interval_min
+ new_start = sleep_start_hour if sleep_start_hour is not None else target.sleep_start_hour
+ new_end = sleep_end_hour if sleep_end_hour is not None else target.sleep_end_hour
+ new_enabled = is_sleep_enabled if is_sleep_enabled is not None else target.is_sleep_enabled
+
+ await conn.execute(
+ """
+ UPDATE targets
+ SET post_interval_min = $1, sleep_start_hour = $2, sleep_end_hour = $3, is_sleep_enabled = $4
+ WHERE id = $5;
+ """,
+ new_interval, new_start, new_end, new_enabled, target_id
+ )
+
async def get_active_targets(self) -> List[TargetChannel]:
pool = await self._get_pool()
async with pool.acquire() as conn:
diff --git a/main.py b/main.py
index 066161e..ed47dc4 100644
--- a/main.py
+++ b/main.py
@@ -9,7 +9,6 @@ from db.repository import Repository
from core.metrics import start_metrics_server
from core.llm import LLMClient
from core.queue import RedisQueue
-from services.queue_consumer import QueueConsumerService
from services.ai_processor import AIProcessor
from services.collector import CollectorService
from services.admin_bot import AdminBotService
@@ -26,7 +25,7 @@ logging.basicConfig(
logger = logging.getLogger("copykar.main")
async def main():
- logger.info("Starting Copykar System with Persian Target Rewriting & Review Pipeline...")
+ logger.info("Starting Copykar System with Immediate Admin Review & Paced Target Queues...")
# 1. Start Prometheus metrics server
metrics_port = int(os.getenv("METRICS_PORT", "8000"))
@@ -44,15 +43,17 @@ async def main():
# 3. Create Services
ai_processor = AIProcessor(repo=repo, llm=llm)
- admin_bot = AdminBotService(repo=repo, ai_processor=ai_processor)
- collector = CollectorService(repo=repo, queue=redis_queue, on_post_received=admin_bot.send_raw_review_post)
+ admin_bot = AdminBotService(repo=repo, ai_processor=ai_processor, queue=redis_queue)
+ collector = CollectorService(repo=repo, on_post_received=admin_bot.send_raw_review_post)
admin_bot.set_collector(collector)
- queue_consumer = QueueConsumerService(queue=redis_queue, on_post_popped=admin_bot.send_raw_review_post)
+
+ # Publisher handles per-target delivery queues, intervals, and sleep windows
+ publisher = PublisherService(repo=repo, queue=redis_queue, client=collector.client)
# 4. Start all services
await admin_bot.start()
await collector.start(notify_fn=admin_bot.notify_admins)
- await queue_consumer.start()
+ await publisher.start()
logger.info("All Copykar services are active and running.")
@@ -73,7 +74,7 @@ async def main():
finally:
logger.info("Shutting down Copykar services...")
await collector.stop()
- await queue_consumer.stop()
+ await publisher.stop()
await admin_bot.stop()
await redis_queue.close()
await close_db_pool()
diff --git a/services/admin_bot.py b/services/admin_bot.py
index a0861da..4324bd7 100644
--- a/services/admin_bot.py
+++ b/services/admin_bot.py
@@ -5,6 +5,7 @@ from typing import Optional, List, Dict
from telethon import TelegramClient, events, Button
from db.models import Post, TargetChannel, SourceChannel
from db.repository import Repository
+from core.queue import RedisQueue
from core.metrics import ADMIN_ACTIONS_TOTAL, TARGET_ACTIVITY_TOTAL
from core.proxy import get_telegram_proxy
@@ -17,7 +18,8 @@ def get_persian_main_menu():
[Button.text("📊 آمار و وضعیت ناوگان", resize=True), Button.text("🔑 درخواست کد لاگین", resize=True)],
[Button.text("📡 کانالهای مبدا", resize=True), Button.text("🎯 کانالهای مقصد", resize=True)],
[Button.text("➕ افزودن کانال مبدا", resize=True), Button.text("➕ افزودن کانال مقصد", resize=True)],
- [Button.text("🎭 تنظیم شخصیت کانالها", resize=True), Button.text("❓ راهنمای سیستم", resize=True)]
+ [Button.text("🎭 تنظیم شخصیت کانالها", resize=True), Button.text("⏰ زمانبندی و خواب کانالها", resize=True)],
+ [Button.text("❓ راهنمای سیستم", resize=True)]
]
class AdminBotService:
@@ -25,6 +27,7 @@ class AdminBotService:
self,
repo: Repository,
ai_processor = None,
+ queue: Optional[RedisQueue] = None,
bot_token: Optional[str] = None,
api_id: Optional[int] = None,
api_hash: Optional[str] = None,
@@ -34,6 +37,7 @@ class AdminBotService:
):
self.repo = repo
self.ai_processor = ai_processor
+ self.queue = queue
self.bot_token = bot_token or os.getenv("BOT_TOKEN", "")
self.api_id = api_id or int(os.getenv("API_ID", "0"))
self.api_hash = api_hash or os.getenv("API_HASH", "")
@@ -44,7 +48,6 @@ class AdminBotService:
os.makedirs(os.path.dirname(self.session_name), exist_ok=True)
self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy())
self.collector = None
- # In-memory store for pending rewritten previews: {f"{post_id}:{target_id}": rewritten_text}
self.preview_cache: Dict[str, str] = {}
def set_collector(self, collector):
@@ -53,6 +56,9 @@ class AdminBotService:
def set_ai_processor(self, ai_processor):
self.ai_processor = ai_processor
+ def set_queue(self, queue: RedisQueue):
+ self.queue = queue
+
def is_admin(self, user_id: int) -> bool:
return not self.admin_user_ids or user_id in self.admin_user_ids
@@ -78,7 +84,6 @@ class AdminBotService:
def _build_raw_post_keyboard(self, post: Post, targets: List[TargetChannel]):
buttons = []
- # Check which targets this post has already been sent to
sent_target_ids = set()
if post.published_to:
for item in post.published_to:
@@ -93,6 +98,9 @@ class AdminBotService:
if len(row) == 2:
buttons.append(row)
row = []
+ if row:
+ buttons.append(row)
+
buttons.append([
Button.inline("❌ رد و بایگانی", data=f"rej:{post.id}"),
Button.inline("🗑 حذف پیام از کانال", data=f"del_msg:{post.id}")
@@ -102,7 +110,7 @@ class AdminBotService:
def _format_raw_post_caption(self, post: Post) -> str:
published_lines = ""
if post.published_to:
- published_lines = "📤 ارسال شده به کانالهای:\n"
+ published_lines = "📤 ارسال شده / در صف ارسال کانالهای:\n"
for item in post.published_to:
if isinstance(item, dict):
t_title = item.get("target_title", "کانال مقصد")
@@ -119,7 +127,7 @@ class AdminBotService:
async def send_raw_review_post(self, post_id: int):
post = await self.repo.get_post_by_id(post_id)
- if not post or not self.review_channel_id:
+ if not post or not self.review_channel_id or post.is_deleted:
return
targets = await self.repo.get_active_targets()
@@ -163,7 +171,7 @@ class AdminBotService:
"👋 به پنل مدیریت سیستم هوشمند کپیکار خوش آمدید!\n\n"
f"• 🤖 وضعیت ربات جمعآوریکننده: {userbot_status}\n"
f"• 📋 شناسه کانال ادمینها: {self.review_channel_id}\n\n"
- "از دکمههای زیر برای مدیریت کانالها، تنظیم شخصیت و آمار استفاده کنید:"
+ "از دکمههای زیر برای مدیریت کانالها، تنظیم شخصیت و زمانبندی استفاده کنید:"
)
await event.reply(welcome_text, parse_mode="html", buttons=get_persian_main_menu())
@@ -175,11 +183,11 @@ class AdminBotService:
pending_review = len(await self.repo.get_posts_by_status("pending_review", limit=5000))
published = len(await self.repo.get_posts_by_status("published", limit=5000))
rejected = len(await self.repo.get_posts_by_status("rejected", limit=5000))
- redis_q = await self.collector.queue.qsize() if (self.collector and self.collector.queue) else 0
+ total_redis_q = await self.queue.get_total_queued_posts() if self.queue else 0
text = (
"📊 آمار زنده سیستم کپیکار:\n\n"
- f"• 📥 پستهای موجود در صف ردیس: {redis_q}\n"
+ f"• 📥 مجموع پستهای در صف ارسال کانالهای مقصد: {total_redis_q}\n"
f"• 📋 پستهای در انتظار بررسی ادمین: {pending_review}\n"
f"• 🚀 پستهای منتشر شده: {published}\n"
f"• ❌ پستهای رد شده: {rejected}\n\n"
@@ -300,10 +308,15 @@ class AdminBotService:
return
lines = ["🎯 کانالهای مقصد برای انتشار:\n"]
for t in targets:
+ qsize = await self.queue.get_target_queue_size(t.id) if self.queue else 0
+ sleep_info = f"{t.sleep_start_hour}:00 تا {t.sleep_end_hour}:00" if t.is_sleep_enabled else "غیرفعال"
lines.append(
f"• {t.title} (شناسه: {t.channel_id} | ID دیتابیس: {t.id})\n"
- f" 🎭 شخصیت و لحن: {t.personality or 'پیشفرض'}\n"
- f" 🏷 فوتر / تگها: {t.custom_footer or 'ندارد'}\n"
+ f" ⏱ فاصله ارسال: هر {t.post_interval_min} دقیقه\n"
+ f" 🌙 ساعت خواب: {sleep_info}\n"
+ f" 📥 تعداد در صف ارسال: {qsize} پست\n"
+ f" 🎭 شخصیت: {t.personality or 'پیشفرض'}\n"
+ f" 🏷 فوتر: {t.custom_footer or 'ندارد'}\n"
)
await event.reply("\n".join(lines), parse_mode="html", buttons=get_persian_main_menu())
@@ -330,11 +343,94 @@ class AdminBotService:
tid = await self.repo.add_target(channel_id=ch_id, title=title, username=username)
await event.reply(
f"✅ کانال مقصد {title} افزوده شد (ID دیتابیس: {tid}).\n\n"
- f"اکنون میتوانید با دکمه 🎭 تنظیم شخصیت کانالها لحن و تگهای آن را تنظیم کنید.",
+ f"اکنون میتوانید با دکمه 🎭 تنظیم شخصیت کانالها لحن و با ⏰ زمانبندی و خواب فواصل ارسال را تنظیم کنید.",
parse_mode="html",
buttons=get_persian_main_menu()
)
+ # --- Schedule & Sleep Management ---
+ @self.client.on(events.NewMessage(pattern=r"(?i)^(/schedule|⏰ زمانبندی و خواب کانالها)$"))
+ async def cmd_schedule(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ targets = await self.repo.get_active_targets()
+ if not targets:
+ await event.reply("ابتدا یک کانال مقصد اضافه کنید.", parse_mode="html")
+ return
+
+ text = (
+ "⏰ مدیریت فاصله ارسال و ساعت خواب کانالها:\n\n"
+ "1. تنظیم فاصله ارسال پستها (به دقیقه):\n"
+ "/set_interval <شناسه_دیتابیس_کانال> <دقیقه>\n"
+ "مثال (ارسال هر ۲۰ دقیقه یک پست):\n"
+ "/set_interval 1 20\n\n"
+ "2. تنظیم ساعت خواب (عدم ارسال پیام در این ساعات):\n"
+ "/set_sleep <شناسه_دیتابیس_کانال> <ساعت_شروع> <ساعت_پایان>\n"
+ "مثال (خواب از ساعت ۲۳ شب تا ۸ صبح):\n"
+ "/set_sleep 1 23 8\n\n"
+ "3. غیرفعال کردن ساعت خواب:\n"
+ "/disable_sleep <شناسه_دیتابیس_کانال>\n\n"
+ "وضعیت فعلی کانالها:\n"
+ )
+ for t in targets:
+ qsize = await self.queue.get_target_queue_size(t.id) if self.queue else 0
+ sleep_st = f"🌙 خواب از {t.sleep_start_hour}:00 تا {t.sleep_end_hour}:00" if t.is_sleep_enabled else "☀️ بدون ساعت خواب"
+ text += f"• ID: {t.id} | {t.title} ➔ هر {t.post_interval_min} دقیقه | {sleep_st} (صف: {qsize} پست)\n"
+
+ await event.reply(text, parse_mode="html", buttons=get_persian_main_menu())
+
+ @self.client.on(events.NewMessage(pattern=r"^/set_interval\s+(\d+)\s+(\d+)"))
+ async def cmd_set_interval(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ target_id = int(event.pattern_match.group(1))
+ interval_min = int(event.pattern_match.group(2))
+ target = await self.repo.get_target_by_id(target_id)
+ if not target:
+ await event.reply(f"❌ کانال مقصد با شناسه {target_id} یافت نشد.")
+ return
+ await self.repo.update_target_schedule(target_id=target_id, post_interval_min=interval_min)
+ await event.reply(
+ f"✅ فاصله ارسال برای کانال {target.title} به هر {interval_min} دقیقه تنظیم شد.",
+ parse_mode="html",
+ buttons=get_persian_main_menu()
+ )
+
+ @self.client.on(events.NewMessage(pattern=r"^/set_sleep\s+(\d+)\s+(\d+)\s+(\d+)"))
+ async def cmd_set_sleep(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ target_id = int(event.pattern_match.group(1))
+ start_h = int(event.pattern_match.group(2))
+ end_h = int(event.pattern_match.group(3))
+ target = await self.repo.get_target_by_id(target_id)
+ if not target:
+ await event.reply(f"❌ کانال مقصد با شناسه {target_id} یافت نشد.")
+ return
+ await self.repo.update_target_schedule(
+ target_id=target_id,
+ sleep_start_hour=start_h,
+ sleep_end_hour=end_h,
+ is_sleep_enabled=True
+ )
+ await event.reply(
+ f"🌙 ساعت خواب برای کانال {target.title} از ساعت {start_h}:00 تا {end_h}:00 فعال شد.\n(در این بازه هیچ پیامی ارسال نخواهد شد و در صف منتظر میماند).",
+ parse_mode="html",
+ buttons=get_persian_main_menu()
+ )
+
+ @self.client.on(events.NewMessage(pattern=r"^/disable_sleep\s+(\d+)"))
+ async def cmd_disable_sleep(event: events.NewMessage.Event):
+ if not self.is_admin(event.sender_id):
+ return
+ target_id = int(event.pattern_match.group(1))
+ target = await self.repo.get_target_by_id(target_id)
+ if not target:
+ await event.reply(f"❌ کانال مقصد با شناسه {target_id} یافت نشد.")
+ return
+ await self.repo.update_target_schedule(target_id=target_id, is_sleep_enabled=False)
+ await event.reply(f"☀️ ساعت خواب برای کانال {target.title} غیرفعال شد.", parse_mode="html", buttons=get_persian_main_menu())
+
# --- Channel Personality & Tags Configuration ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/personality|🎭 تنظیم شخصیت کانالها)$"))
async def cmd_personality(event: events.NewMessage.Event):
@@ -406,14 +502,13 @@ class AdminBotService:
return
help_text = (
"📖 راهنمای فرآیند کاری سیستم کپیکار:\n\n"
- "1. 📥 دریافت خام پستها: پستها بدون پردازش هوش مصنوعی مستقیماً به کانال ادمینها ارسال میشوند.\n"
- "2. 🎯 انتخاب کانال مقصد: با لمس دکمه هر کانال، هوش مصنوعی پست را متناسب با شخصیت، استایل و فوتر اختصاصی همان کانال بازنویسی کرده و تمام تگها و لینکهای مبدا را حذف میکند.\n"
- "3. 👁 پیشنمایش زنده: پیشنمایش بازنویسی شده به همراه دکمه تایید نهایی نمایش داده میشود.\n"
- "4. 🚀 انتشار و ارسال مجدد: پس از انتشار، پست اصلی در کانال ادمین بازگردانده شده و سابقه انتشار نمایش مییابد تا بتوانید آن را به سایر کانالها نیز ارسال کنید."
+ "1. 📥 دریافت آنی: پستهای مبدا فوری و بدون تاخیر در کانال ادمین قرار میگیرند.\n"
+ "2. 🎯 بازنویسی بر اساس مقصد: با زدن دکمه کانال مقصد، هوش مصنوعی پست را متناسب با شخصیت و فوتر آن کانال بازنویسی میکند.\n"
+ "3. ⏳ صف ارسال زمانبندی شده: پس از تایید، پست در صف ردیس کانال مقصد قرار میگیرد و با رعایت فاصله زمانی (Interval) و ساعات خواب (Sleep) به ترتیب منتشر میشود."
)
await event.reply(help_text, parse_mode="html", buttons=get_persian_main_menu())
- # --- Interactive Inline Callbacks for Reviews & Target Rewrites ---
+ # --- Inline Callback Queries ---
@self.client.on(events.CallbackQuery)
async def on_callback(event: events.CallbackQuery.Event):
if not self.is_admin(event.sender_id):
@@ -454,7 +549,6 @@ class AdminBotService:
await event.answer(f"در حال بازنویسی برای {target.title}...")
- # Show loading placeholder
loading_caption = (
f"🤖 در حال بازنویسی هوشمند برای کانال: {target.title}...\n"
f"(اعمال لحن اختصاصی و حذف تگهای مبدا)"
@@ -464,24 +558,22 @@ class AdminBotService:
except Exception:
pass
- # Run AI Rewrite
rewritten_text = await self.ai_processor.rewrite_for_target(post.raw_text or "", target)
cache_key = f"{post_id}:{target_id}"
self.preview_cache[cache_key] = rewritten_text
- # Build Preview Card
preview_caption = (
f"🎯 پیشنمایش بازنویسی شده برای: {target.title}\n"
f"🎭 شخصیت و لحن: {target.personality or 'پیشفرض'}\n"
f"➖➖➖➖➖➖➖➖➖➖\n\n"
f"{rewritten_text}\n\n"
f"➖➖➖➖➖➖➖➖➖➖\n"
- f"آیا این متن مورد تایید است؟"
+ f"آیا این متن برای صف انتشار تایید است؟"
)
preview_buttons = [
[
- Button.inline(f"✅ تایید و ارسال به {target.title}", data=f"pub:{post_id}:{target_id}"),
+ Button.inline(f"✅ تایید و افزودن به صف {target.title}", data=f"pub:{post_id}:{target_id}"),
],
[
Button.inline("🔙 انصراف / بازگشت به پست اصلی", data=f"cancel:{post_id}")
@@ -490,7 +582,7 @@ class AdminBotService:
await event.edit(preview_caption, parse_mode="html", buttons=preview_buttons)
- # 3. Publish to Target Confirmed
+ # 3. Add to Target's Redis Queue
elif data.startswith("pub:"):
_, post_id_str, target_id_str = data.split(":")
post_id = int(post_id_str)
@@ -505,45 +597,33 @@ class AdminBotService:
cache_key = f"{post_id}:{target_id}"
text_to_publish = self.preview_cache.get(cache_key) or post.raw_text or ""
- await event.answer(f"در حال ارسال به {target.title}...")
+ # Push to Target's Redis Queue
+ payload = {
+ "post_id": post.id,
+ "text": text_to_publish,
+ "media_path": post.media_path,
+ "target_id": target.id,
+ "target_title": target.title
+ }
+ if self.queue:
+ await self.queue.push_target_post(target.id, payload)
+ ADMIN_ACTIONS_TOTAL.labels(action="approved").inc()
- try:
- # Publish via userbot or bot
- client_to_use = self.collector.client if (self.collector and self.collector.client.is_connected()) else self.client
- if post.media_path and os.path.exists(post.media_path):
- await client_to_use.send_file(
- target.channel_id,
- file=post.media_path,
- caption=text_to_publish,
- parse_mode="html"
- )
- else:
- await client_to_use.send_message(
- target.channel_id,
- text_to_publish,
- parse_mode="html"
- )
+ # Update database record
+ await self.repo.record_post_published_to_target(post_id, target.id, target.title or "Target")
- # Record publication in database
- await self.repo.record_post_published_to_target(post_id, target.id, target.title or "Target")
- await self.repo.update_target_last_post(target.id)
- ADMIN_ACTIONS_TOTAL.labels(action="approved").inc()
- TARGET_ACTIVITY_TOTAL.labels(channel_id=str(target.channel_id), title=target.title or '').inc()
+ # Reload updated post with publication history
+ updated_post = await self.repo.get_post_by_id(post_id)
+ targets = await self.repo.get_active_targets()
+ new_caption = self._format_raw_post_caption(updated_post)
+ new_buttons = self._build_raw_post_keyboard(updated_post, targets)
- # Reload updated post with publication history
- updated_post = await self.repo.get_post_by_id(post_id)
- targets = await self.repo.get_active_targets()
- new_caption = self._format_raw_post_caption(updated_post)
- new_buttons = self._build_raw_post_keyboard(updated_post, targets)
-
- await event.edit(
- f"✅ با موفقیت در {target.title} منتشر شد!\n\n{new_caption}",
- parse_mode="html",
- buttons=new_buttons
- )
- except Exception as e:
- logger.error(f"Failed to publish to {target.channel_id}: {e}", exc_info=True)
- await event.answer(f"❌ خطا در ارسال به کانال: {e}", alert=True)
+ await event.edit(
+ f"✅ به صف انتشار کانال {target.title} اضافه شد!\n\n{new_caption}",
+ parse_mode="html",
+ buttons=new_buttons
+ )
+ await event.answer(f"به صف {target.title} افزوده شد!")
# 4. Cancel Preview & Restore Original Card
elif data.startswith("cancel:"):
diff --git a/services/collector.py b/services/collector.py
index e6c21b5..1d155e1 100644
--- a/services/collector.py
+++ b/services/collector.py
@@ -143,9 +143,7 @@ class CollectorService:
SOURCE_ACTIVITY_TOTAL.labels(channel_id=str(chat_id), title=source.title or 'Unknown').inc()
logger.info(f"Collected raw post ID {post_id} from source channel {chat_id}")
- if self.queue:
- await self.queue.push(post_id)
- elif self.on_post_received:
+ if self.on_post_received:
await self.on_post_received(post_id)
except Exception as e:
logger.error(f"Error handling message from {event.chat_id}: {e}", exc_info=True)
@@ -216,9 +214,7 @@ class CollectorService:
SOURCE_ACTIVITY_TOTAL.labels(channel_id=str(channel_id), title=source_title).inc()
logger.info(f"Backfilled raw post ID {post_id} from {channel_id}")
- if self.queue:
- await self.queue.push(post_id)
- elif self.on_post_received:
+ if self.on_post_received:
await self.on_post_received(post_id)
else:
skipped_count += 1
diff --git a/services/publisher.py b/services/publisher.py
index f390190..3382444 100644
--- a/services/publisher.py
+++ b/services/publisher.py
@@ -2,10 +2,12 @@ import os
import asyncio
import logging
from datetime import datetime, timezone
-from typing import Optional
+from typing import Optional, List
from telethon import TelegramClient
+from db.models import TargetChannel
from db.repository import Repository
-from core.metrics import POSTS_PUBLISHED_TOTAL, QUEUE_POSTS_GAUGE
+from core.queue import RedisQueue
+from core.metrics import TARGET_ACTIVITY_TOTAL, QUEUE_POSTS_GAUGE, REDIS_QUEUE_SIZE_GAUGE
from core.proxy import get_telegram_proxy
logger = logging.getLogger(__name__)
@@ -16,48 +18,79 @@ class PublisherService:
def __init__(
self,
repo: Repository,
+ queue: RedisQueue,
+ client: Optional[TelegramClient] = None,
api_id: Optional[int] = None,
api_hash: Optional[str] = None,
- session_name: Optional[str] = None,
bot_token: Optional[str] = None,
+ session_name: Optional[str] = None,
):
self.repo = repo
+ self.queue = queue
+ self.client = client
self.api_id = api_id or int(os.getenv("API_ID", "0"))
self.api_hash = api_hash or os.getenv("API_HASH", "")
self.bot_token = bot_token or os.getenv("BOT_TOKEN")
self.session_name = session_name or os.path.join(SESSION_DIR, "publisher.session")
os.makedirs(os.path.dirname(self.session_name), exist_ok=True)
- self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy())
+ if not self.client:
+ self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy())
self._running = False
self._task: Optional[asyncio.Task] = None
async def start(self):
- logger.info("Starting Publisher Service...")
- if self.bot_token:
- await self.client.start(bot_token=self.bot_token)
- else:
- await self.client.start()
- logger.info("Publisher Service connected successfully.")
-
+ logger.info("Starting Paced Target Publisher Service...")
+ if not self.client.is_connected():
+ if self.bot_token:
+ await self.client.start(bot_token=self.bot_token)
+ else:
+ await self.client.start()
+ logger.info("Paced Target Publisher Service connected.")
self._running = True
self._task = asyncio.create_task(self._publisher_loop())
+ def _is_in_sleep_window(self, target: TargetChannel, current_hour: int) -> bool:
+ if not target.is_sleep_enabled:
+ return False
+ start = target.sleep_start_hour
+ end = target.sleep_end_hour
+ if start == end:
+ return False
+ if start < end:
+ return start <= current_hour < end
+ else:
+ # Overnight sleep (e.g. 23:00 to 08:00)
+ return current_hour >= start or current_hour < end
+
async def _publisher_loop(self):
while self._running:
try:
- await self._process_pending_queues()
+ await self._process_all_target_queues()
except Exception as e:
- logger.error(f"Error in publisher loop: {e}", exc_info=True)
+ logger.error(f"Error in target publisher loop: {e}", exc_info=True)
await asyncio.sleep(15)
- async def _process_pending_queues(self):
+ async def _process_all_target_queues(self):
targets = await self.repo.get_active_targets()
now = datetime.now(timezone.utc)
+ current_hour_local = (now.hour + 3) % 24 # UTC+3:30 approx hour
- pending_count = len(await self.repo.get_posts_by_status("approved", limit=5000))
- QUEUE_POSTS_GAUGE.labels(status="approved").set(pending_count)
+ total_queued = 0
for target in targets:
+ qsize = await self.queue.get_target_queue_size(target.id)
+ total_queued += qsize
+ QUEUE_POSTS_GAUGE.labels(status=f"target_{target.id}").set(qsize)
+
+ if qsize == 0:
+ continue
+
+ # 1. Check Sleep Window
+ if self._is_in_sleep_window(target, current_hour_local):
+ logger.debug(f"Target #{target.id} ({target.title}) in sleep window ({target.sleep_start_hour}:00-{target.sleep_end_hour}:00). Skipping.")
+ continue
+
+ # 2. Check Cooldown Interval
if target.last_post_time:
last_post = target.last_post_time
if last_post.tzinfo is None:
@@ -66,37 +99,42 @@ class PublisherService:
if diff_minutes < target.post_interval_min:
continue
- post = await self.repo.get_next_approved_post_for_target(target.id)
- if not post:
+ # 3. Pop next post payload for this target
+ payload = await self.queue.pop_target_post(target.id)
+ if not payload:
continue
+ post_id = payload.get("post_id")
+ text = payload.get("text", "")
+ media_path = payload.get("media_path")
+
try:
- publish_text = post.ai_text or post.raw_text or ""
- if post.media_path and os.path.exists(post.media_path):
+ if media_path and os.path.exists(media_path):
await self.client.send_file(
target.channel_id,
- file=post.media_path,
- caption=publish_text,
- parse_mode="markdown"
+ file=media_path,
+ caption=text,
+ parse_mode="html"
)
else:
await self.client.send_message(
target.channel_id,
- publish_text,
- parse_mode="markdown"
+ text,
+ parse_mode="html"
)
- await self.repo.mark_post_published(post.id)
+ # Record publication in database
+ await self.repo.record_post_published_to_target(post_id, target.id, target.title or "Target")
await self.repo.update_target_last_post(target.id)
- POSTS_PUBLISHED_TOTAL.labels(target_channel_id=str(target.channel_id)).inc()
- logger.info(f"Successfully published post {post.id} to target channel {target.title} ({target.channel_id})")
+ TARGET_ACTIVITY_TOTAL.labels(channel_id=str(target.channel_id), title=target.title or '').inc()
+ logger.info(f"Published post ID {post_id} to Target {target.title} ({target.channel_id})")
except Exception as e:
- logger.error(f"Failed to publish post {post.id} to target {target.channel_id}: {e}", exc_info=True)
+ logger.error(f"Failed to publish queued post {post_id} to target {target.channel_id}: {e}", exc_info=True)
+
+ REDIS_QUEUE_SIZE_GAUGE.set(total_queued)
async def stop(self):
self._running = False
if self._task:
self._task.cancel()
- if self.client.is_connected():
- await self.client.disconnect()
- logger.info("Publisher Service disconnected.")
+ logger.info("Publisher Service stopped.")