feat(workflow): implement target channel personalities, raw review cards, on-demand AI rewrites and multi-target dispatch

This commit is contained in:
mamad
2026-08-27 22:11:33 +03:30
parent 4f5570850e
commit 26cd91a48f
10 changed files with 695 additions and 699 deletions
+14 -2
View File
@@ -10,6 +10,18 @@ COLLECTED_POSTS_TOTAL = Counter(
["source_channel_id"]
)
SOURCE_ACTIVITY_TOTAL = Counter(
"copykar_source_activity_total",
"Total posts ingested per source channel",
["channel_id", "title"]
)
TARGET_ACTIVITY_TOTAL = Counter(
"copykar_target_activity_total",
"Total posts published per target channel",
["channel_id", "title"]
)
AI_REQUESTS_TOTAL = Counter(
"copykar_ai_requests_total",
"Total AI API calls made",
@@ -19,13 +31,13 @@ AI_REQUESTS_TOTAL = Counter(
DUPLICATES_DETECTED_TOTAL = Counter(
"copykar_duplicates_detected_total",
"Total duplicate posts detected",
["method"] # "hash" or "ai_semantic"
["method"]
)
ADMIN_ACTIONS_TOTAL = Counter(
"copykar_admin_actions_total",
"Total review decisions by admins",
["action"] # "approved", "rejected", "routed"
["action"]
)
POSTS_PUBLISHED_TOTAL = Counter(
+10 -1
View File
@@ -1,5 +1,6 @@
import asyncpg
import os
import json
from typing import Optional
DATABASE_URL = os.getenv(
@@ -23,6 +24,8 @@ CREATE TABLE IF NOT EXISTS targets (
title VARCHAR(255),
username VARCHAR(255),
post_interval_min INT DEFAULT 30,
personality TEXT DEFAULT '',
custom_footer TEXT DEFAULT '',
last_post_time TIMESTAMPTZ,
is_active BOOLEAN DEFAULT TRUE,
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
@@ -44,7 +47,8 @@ CREATE TABLE IF NOT EXISTS posts (
ai_text TEXT,
suggested_target_id INT REFERENCES targets(id) ON DELETE SET NULL,
target_channel_id INT REFERENCES targets(id) ON DELETE SET NULL,
status VARCHAR(32) DEFAULT 'pending_ai',
status VARCHAR(32) DEFAULT 'pending_review',
published_to JSONB DEFAULT '[]'::jsonb,
review_message_id BIGINT,
scheduled_at TIMESTAMPTZ,
published_at TIMESTAMPTZ,
@@ -63,6 +67,11 @@ CREATE TABLE IF NOT EXISTS settings (
value TEXT NOT NULL,
description TEXT
);
-- 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 posts ADD COLUMN IF NOT EXISTS published_to JSONB DEFAULT '[]'::jsonb;
"""
_pool: Optional[asyncpg.Pool] = None
+5 -2
View File
@@ -1,5 +1,5 @@
from dataclasses import dataclass, field
from typing import Optional, List
from typing import Optional, List, Dict, Any
@dataclass
class SourceChannel:
@@ -17,6 +17,8 @@ class TargetChannel:
title: Optional[str]
username: Optional[str]
post_interval_min: int = 30
personality: str = ""
custom_footer: str = ""
last_post_time: Optional[str] = None
is_active: bool = True
created_at: Optional[str] = None
@@ -38,7 +40,8 @@ class Post:
ai_text: Optional[str] = None
suggested_target_id: Optional[int] = None
target_channel_id: Optional[int] = None
status: str = "pending_ai" # pending_ai, pending_review, approved, scheduled, published, rejected
status: str = "pending_review" # pending_review, published, rejected
published_to: List[Dict[str, Any]] = field(default_factory=list)
review_message_id: Optional[int] = None
scheduled_at: Optional[str] = None
published_at: Optional[str] = None
+69 -139
View File
@@ -1,8 +1,20 @@
import json
import asyncpg
from typing import List, Optional
from typing import List, Optional, Dict, Any
from db.models import SourceChannel, TargetChannel, Post, Setting
from db.database import get_db_pool
def _parse_post_row(row: asyncpg.Record) -> Post:
data = dict(row)
if isinstance(data.get("published_to"), str):
try:
data["published_to"] = json.loads(data["published_to"])
except Exception:
data["published_to"] = []
elif data.get("published_to") is None:
data["published_to"] = []
return Post(**data)
class Repository:
def __init__(self, dsn: Optional[str] = None):
self.dsn = dsn
@@ -33,7 +45,7 @@ class Repository:
async def get_active_sources(self) -> List[SourceChannel]:
pool = await self._get_pool()
async with pool.acquire() as conn:
rows = await conn.fetch("SELECT * FROM sources WHERE is_active = TRUE;")
rows = await conn.fetch("SELECT * FROM sources WHERE is_active = TRUE ORDER BY id ASC;")
return [SourceChannel(**dict(r)) for r in rows]
async def get_source_by_channel_id(self, channel_id: int) -> Optional[SourceChannel]:
@@ -43,13 +55,21 @@ class Repository:
return SourceChannel(**dict(row)) if row else None
# --- Target Channels ---
async def add_target(self, channel_id: int, title: Optional[str] = None, username: Optional[str] = None, post_interval_min: int = 30) -> int:
async def add_target(
self,
channel_id: int,
title: Optional[str] = None,
username: Optional[str] = None,
post_interval_min: int = 30,
personality: str = "",
custom_footer: str = ""
) -> int:
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow(
"""
INSERT INTO targets (channel_id, title, username, post_interval_min)
VALUES ($1, $2, $3, $4)
INSERT INTO targets (channel_id, title, username, post_interval_min, personality, custom_footer)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT(channel_id) DO UPDATE SET
title = EXCLUDED.title,
username = EXCLUDED.username,
@@ -57,14 +77,36 @@ class Repository:
is_active = TRUE
RETURNING id;
""",
channel_id, title, username, post_interval_min,
channel_id, title, username, post_interval_min, personality, custom_footer,
)
return row["id"]
async def update_target_personality(self, target_id: int, personality: str, custom_footer: Optional[str] = None) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
if custom_footer is not None:
await conn.execute(
"UPDATE targets SET personality = $1, custom_footer = $2 WHERE id = $3;",
personality, custom_footer, target_id
)
else:
await conn.execute(
"UPDATE targets SET personality = $1 WHERE id = $2;",
personality, target_id
)
async def update_target_footer(self, target_id: int, custom_footer: str) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"UPDATE targets SET custom_footer = $1 WHERE id = $2;",
custom_footer, target_id
)
async def get_active_targets(self) -> List[TargetChannel]:
pool = await self._get_pool()
async with pool.acquire() as conn:
rows = await conn.fetch("SELECT * FROM targets WHERE is_active = TRUE;")
rows = await conn.fetch("SELECT * FROM targets WHERE is_active = TRUE ORDER BY id ASC;")
return [TargetChannel(**dict(r)) for r in rows]
async def get_target_by_id(self, target_id: int) -> Optional[TargetChannel]:
@@ -81,34 +123,7 @@ class Repository:
target_id,
)
# --- Posts & Deduplication ---
async def find_duplicate_post_by_hash(self, content_hash: str) -> Optional[Post]:
if not content_hash:
return None
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow(
"SELECT * FROM posts WHERE content_hash = $1 ORDER BY id ASC LIMIT 1;",
content_hash,
)
return Post(**dict(row)) if row else None
async def find_candidate_posts_by_tags(self, tags: List[str], exclude_post_id: Optional[int] = None, hours_lookback: int = 72, limit: int = 5) -> List[Post]:
if not tags:
return []
pool = await self._get_pool()
async with pool.acquire() as conn:
query = """
SELECT * FROM posts
WHERE tags && $1::text[]
AND created_at >= NOW() - ($2 || ' hours')::interval
AND ($3::bigint IS NULL OR id != $3::bigint)
ORDER BY created_at DESC
LIMIT $4;
"""
rows = await conn.fetch(query, tags, str(hours_lookback), exclude_post_id, limit)
return [Post(**dict(r)) for r in rows]
# --- Posts & Multi-Channel Dispatch ---
async def create_raw_post(
self,
source_channel_id: int,
@@ -130,7 +145,7 @@ class Repository:
source_channel_id, source_message_id, raw_text, media_path,
media_type, content_hash, is_duplicate, duplicate_of_id, similarity_reason, status
)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, 'pending_ai')
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, 'pending_review')
RETURNING id;
""",
source_channel_id,
@@ -147,6 +162,12 @@ class Repository:
except asyncpg.UniqueViolationError:
return None
async def get_post_by_id(self, post_id: int) -> Optional[Post]:
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow("SELECT * FROM posts WHERE id = $1;", post_id)
return _parse_post_row(row) if row else None
async def get_posts_by_status(self, status: str, limit: int = 20) -> List[Post]:
pool = await self._get_pool()
async with pool.acquire() as conn:
@@ -154,63 +175,7 @@ class Repository:
"SELECT * FROM posts WHERE status = $1 ORDER BY id ASC LIMIT $2;",
status, limit,
)
return [Post(**dict(r)) for r in rows]
async def get_post_by_id(self, post_id: int) -> Optional[Post]:
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow("SELECT * FROM posts WHERE id = $1;", post_id)
return Post(**dict(row)) if row else None
async def update_post_tags(self, post_id: int, tags: List[str], subject: str) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"UPDATE posts SET tags = $1, subject = $2 WHERE id = $3;",
tags, subject, post_id,
)
async def update_post_duplicate_status(
self,
post_id: int,
is_duplicate: bool,
duplicate_of_id: Optional[int] = None,
similarity_reason: Optional[str] = None,
) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"""
UPDATE posts
SET is_duplicate = $1, duplicate_of_id = $2, similarity_reason = $3
WHERE id = $4;
""",
is_duplicate, duplicate_of_id, similarity_reason, post_id,
)
async def update_ai_result(
self,
post_id: int,
subject: str,
ai_text: str,
tags: List[str],
suggested_target_id: Optional[int] = None,
is_duplicate: bool = False,
duplicate_of_id: Optional[int] = None,
similarity_reason: Optional[str] = None,
) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"""
UPDATE posts
SET subject = $1, ai_text = $2, tags = $3, suggested_target_id = $4,
is_duplicate = $5, duplicate_of_id = $6, similarity_reason = $7,
status = 'pending_review'
WHERE id = $8;
""",
subject, ai_text, tags, suggested_target_id, is_duplicate, duplicate_of_id, similarity_reason, post_id,
)
return [_parse_post_row(r) for r in rows]
async def update_review_message_id(self, post_id: int, review_message_id: int) -> None:
pool = await self._get_pool()
@@ -220,62 +185,27 @@ class Repository:
review_message_id, post_id,
)
async def approve_post(self, post_id: int, target_channel_id: int) -> None:
async def record_post_published_to_target(self, post_id: int, target_id: int, target_title: str) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
record_item = json.dumps({
"target_id": target_id,
"target_title": target_title,
"published_at": str(asyncpg.types.Type)
})
await conn.execute(
"""
UPDATE posts
SET target_channel_id = $1, status = 'approved'
SET published_to = published_to || $1::jsonb,
status = 'published',
published_at = CURRENT_TIMESTAMP
WHERE id = $2;
""",
target_channel_id, post_id,
f'[{{"target_id": {target_id}, "target_title": "{target_title}", "published_at": "{asyncpg.types.Type}"}}]',
post_id,
)
async def reject_post(self, post_id: int) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute("UPDATE posts SET status = 'rejected' WHERE id = $1;", post_id)
async def mark_post_published(self, post_id: int) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"UPDATE posts SET status = 'published', published_at = CURRENT_TIMESTAMP WHERE id = $1;",
post_id,
)
async def get_next_approved_post_for_target(self, target_id: int) -> Optional[Post]:
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow(
"""
SELECT * FROM posts
WHERE target_channel_id = $1 AND status = 'approved'
ORDER BY id ASC
LIMIT 1;
""",
target_id,
)
return Post(**dict(row)) if row else None
# --- Settings ---
async def get_setting(self, key: str, default: Optional[str] = None) -> Optional[str]:
pool = await self._get_pool()
async with pool.acquire() as conn:
row = await conn.fetchrow("SELECT value FROM settings WHERE key = $1;", key)
return row["value"] if row else default
async def set_setting(self, key: str, value: str, description: Optional[str] = None) -> None:
pool = await self._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"""
INSERT INTO settings (key, value, description)
VALUES ($1, $2, $3)
ON CONFLICT(key) DO UPDATE SET
value = EXCLUDED.value,
description = COALESCE(EXCLUDED.description, settings.description);
""",
key, value, description,
)
+5 -18
View File
@@ -26,7 +26,7 @@ logging.basicConfig(
logger = logging.getLogger("copykar.main")
async def main():
logger.info("Starting Copykar System...")
logger.info("Starting Copykar System with Persian Target Rewriting & Review Pipeline...")
# 1. Start Prometheus metrics server
metrics_port = int(os.getenv("METRICS_PORT", "8000"))
@@ -43,28 +43,16 @@ async def main():
llm = LLMClient()
# 3. Create Services
admin_bot = AdminBotService(repo=repo)
ai_processor = AIProcessor(repo=repo, llm=llm)
# Wrap AI processor to automatically push reviewed posts to admin review channel
original_process_post = ai_processor.process_post
async def process_and_notify(post_id: int):
post = await original_process_post(post_id)
if post:
await admin_bot.send_review_post(post.id)
return post
ai_processor.process_post = process_and_notify
collector = CollectorService(repo=repo, ai_processor=ai_processor, queue=redis_queue)
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.set_collector(collector)
publisher = PublisherService(repo=repo)
queue_consumer = QueueConsumerService(queue=redis_queue, ai_processor=ai_processor)
queue_consumer = QueueConsumerService(queue=redis_queue, on_post_popped=admin_bot.send_raw_review_post)
# 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.")
@@ -76,7 +64,7 @@ async def main():
try:
loop.add_signal_handler(sig, stop_event.set)
except NotImplementedError:
pass # Windows or specific platforms
pass
try:
await stop_event.wait()
@@ -86,7 +74,6 @@ async def main():
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()
+24 -62
View File
@@ -20,7 +20,7 @@
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 0, "y": 1 },
"id": 1,
"title": "Total Collected",
"title": "Total Ingested",
"type": "stat",
"targets": [
{
@@ -66,25 +66,6 @@
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 6, "y": 1 },
"id": 2,
"title": "Duplicates Blocked",
"type": "stat",
"targets": [
{
"expr": "sum(copykar_duplicates_detected_total) or vector(0)",
"legendFormat": "Duplicates",
"refId": "A"
}
],
"fieldConfig": {
"defaults": {
"thresholds": { "mode": "absolute", "steps": [{ "color": "orange", "value": null }] }
}
}
},
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 9, "y": 1 },
"id": 3,
"title": "Admin Approvals",
"type": "stat",
@@ -103,7 +84,7 @@
},
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 12, "y": 1 },
"gridPos": { "h": 4, "w": 3, "x": 9, "y": 1 },
"id": 4,
"title": "Admin Rejections",
"type": "stat",
@@ -122,13 +103,13 @@
},
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 15, "y": 1 },
"gridPos": { "h": 4, "w": 4, "x": 12, "y": 1 },
"id": 5,
"title": "Posts Published",
"title": "Target Deliveries",
"type": "stat",
"targets": [
{
"expr": "sum(copykar_posts_published_total) or vector(0)",
"expr": "sum(copykar_target_activity_total) or vector(0)",
"legendFormat": "Published",
"refId": "A"
}
@@ -141,9 +122,9 @@
},
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 18, "y": 1 },
"gridPos": { "h": 4, "w": 4, "x": 16, "y": 1 },
"id": 6,
"title": "AI Requests Total",
"title": "AI Target Rewrites Total",
"type": "stat",
"targets": [
{
@@ -160,14 +141,14 @@
},
{
"collapsed": false,
"gridPos": { "h": 4, "w": 3, "x": 21, "y": 1 },
"gridPos": { "h": 4, "w": 4, "x": 20, "y": 1 },
"id": 16,
"title": "Publish Queue Depth",
"title": "Review Queue Depth",
"type": "stat",
"targets": [
{
"expr": "copykar_posts_queue_gauge{status=\"approved\"} or vector(0)",
"legendFormat": "Approved Queued",
"expr": "copykar_posts_queue_gauge{status=\"pending_review\"} or vector(0)",
"legendFormat": "Pending Review",
"refId": "A"
}
],
@@ -181,56 +162,37 @@
"collapsed": false,
"gridPos": { "h": 1, "w": 24, "x": 0, "y": 5 },
"id": 101,
"title": "📈 Pipeline Rates & Queue Depths",
"title": "📡 Source & Target Channel Activity Breakdown",
"type": "row"
},
{
"collapsed": false,
"gridPos": { "h": 8, "w": 8, "x": 0, "y": 6 },
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 6 },
"id": 7,
"title": "Ingestion Rate by Source (posts/min)",
"title": "Ingested Posts by Source Channel",
"type": "timeseries",
"targets": [
{
"expr": "rate(copykar_posts_collected_total[1m]) * 60",
"legendFormat": "Source: {{source_channel_id}}",
"expr": "sum by (title) (copykar_source_activity_total)",
"legendFormat": "Source: {{title}}",
"refId": "A"
}
]
},
{
"collapsed": false,
"gridPos": { "h": 8, "w": 8, "x": 8, "y": 6 },
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 6 },
"id": 8,
"title": "Publishing Rate by Target (posts/min)",
"title": "Published Posts by Target Channel",
"type": "timeseries",
"targets": [
{
"expr": "rate(copykar_posts_published_total[1m]) * 60",
"legendFormat": "Target: {{target_channel_id}}",
"expr": "sum by (title) (copykar_target_activity_total)",
"legendFormat": "Target: {{title}}",
"refId": "A"
}
]
},
{
"collapsed": false,
"gridPos": { "h": 8, "w": 8, "x": 16, "y": 6 },
"id": 9,
"title": "Queue Depth Over Time",
"type": "timeseries",
"targets": [
{
"expr": "copykar_redis_queue_size",
"legendFormat": "Redis Incoming Queue (2m Pacing)",
"refId": "A"
},
{
"expr": "copykar_posts_queue_gauge",
"legendFormat": "State: {{status}}",
"refId": "B"
}
]
},
{
"collapsed": false,
"gridPos": { "h": 1, "w": 24, "x": 0, "y": 14 },
@@ -242,12 +204,12 @@
"collapsed": false,
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 15 },
"id": 10,
"title": "Duplicates Intercepted by Method",
"title": "AI Requests by Action & Status",
"type": "timeseries",
"targets": [
{
"expr": "sum by (method) (rate(copykar_duplicates_detected_total[1m]) * 60)",
"legendFormat": "Method: {{method}}",
"expr": "sum by (action, status) (rate(copykar_ai_requests_total[1m]) * 60)",
"legendFormat": "{{action}} ({{status}})",
"refId": "A"
}
]
@@ -338,5 +300,5 @@
"timezone": "browser",
"title": "Copykar Telegram Fleet Executive Dashboard",
"uid": "copykar-executive-dashboard",
"version": 3
"version": 4
}
+491 -330
View File
@@ -1,29 +1,30 @@
import os
import logging
from typing import Optional, List
from datetime import datetime, timezone
from typing import Optional, List, Dict
from telethon import TelegramClient, events, Button
from db.models import Post, TargetChannel
from db.models import Post, TargetChannel, SourceChannel
from db.repository import Repository
from bot.keyboards import get_review_keyboard
from core.metrics import ADMIN_ACTIONS_TOTAL
from core.metrics import ADMIN_ACTIONS_TOTAL, TARGET_ACTIVITY_TOTAL
from core.proxy import get_telegram_proxy
logger = logging.getLogger(__name__)
SESSION_DIR = os.getenv("SESSION_DIR", "/app/sessions" if os.path.exists("/app") else "/projects/telegram-bots/copykar/sessions")
def get_main_menu_keyboard():
def get_persian_main_menu():
return [
[Button.text("🔑 Request Login Code", resize=True), Button.text("📊 Fleet Statistics", resize=True)],
[Button.text("📡 Monitored Sources", resize=True), Button.text("🎯 Target Channels", resize=True)],
[Button.text(" Add Source Guide", resize=True), Button.text(" Add Target Guide", resize=True)],
[Button.text("❓ Help & Documentation", 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:
def __init__(
self,
repo: Repository,
ai_processor = None,
bot_token: Optional[str] = None,
api_id: Optional[int] = None,
api_hash: Optional[str] = None,
@@ -32,6 +33,7 @@ class AdminBotService:
session_name: Optional[str] = None,
):
self.repo = repo
self.ai_processor = ai_processor
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", "")
@@ -42,15 +44,20 @@ 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):
self.collector = collector
def set_ai_processor(self, ai_processor):
self.ai_processor = ai_processor
def is_admin(self, user_id: int) -> bool:
return not self.admin_user_ids or user_id in self.admin_user_ids
async def notify_admins(self, text: str):
"""Broadcast message to review channel and all admin DMs."""
"""Broadcast Persian message to review channel and admin DMs."""
if self.review_channel_id:
try:
await self.client.send_message(self.review_channel_id, text, parse_mode="html")
@@ -59,340 +66,65 @@ class AdminBotService:
for admin_id in self.admin_user_ids:
try:
await self.client.send_message(admin_id, text, parse_mode="html", buttons=get_main_menu_keyboard())
await self.client.send_message(admin_id, text, parse_mode="html", buttons=get_persian_main_menu())
except Exception as e:
logger.debug(f"Could not send DM to admin {admin_id}: {e}")
async def start(self):
logger.info("Starting Admin Review Bot...")
logger.info("Starting Admin Bot Service...")
await self.client.start(bot_token=self.bot_token)
logger.info("Admin Review Bot connected successfully.")
logger.info("Admin Bot connected successfully.")
self._register_handlers()
def _register_handlers(self):
# --- /start and Main Menu ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/start|/menu|menu)$"))
async def cmd_start(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
await event.reply(f"⛔ Unauthorized user ID: <code>{event.sender_id}</code>. Please add this ID to <code>ADMIN_USER_IDS</code> in .env.", parse_mode="html")
return
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:
if isinstance(item, dict) and "target_id" in item:
sent_target_ids.add(int(item["target_id"]))
userbot_status = "🔴 Not Authorized"
if self.collector and self.collector.client.is_connected() and await self.collector.client.is_user_authorized():
me = await self.collector.client.get_me()
userbot_status = f"🟢 Online ({me.first_name})"
row = []
for t in targets:
is_sent = t.id in sent_target_ids
label = f"{t.title}" if is_sent else f"🎯 {t.title}"
row.append(Button.inline(label, data=f"sel_trg:{post.id}:{t.id}"))
if len(row) == 2:
buttons.append(row)
row = []
if row:
buttons.append(row)
welcome_text = (
"👋 <b>Welcome to Copykar Admin Console!</b>\n\n"
f"• 🤖 <b>Userbot Status:</b> {userbot_status}\n"
f"• 📋 <b>Review Channel:</b> <code>{self.review_channel_id}</code>\n\n"
"Use the interactive menu buttons below to manage the fleet:"
)
await event.reply(welcome_text, parse_mode="html", buttons=get_main_menu_keyboard())
buttons.append([Button.inline("❌ رد و بایگانی پست", data=f"rej:{post.id}")])
return buttons
# --- Interactive Userbot Authentication Commands ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/request_code|🔑 Request Login Code)$"))
async def cmd_request_code(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
await event.reply("Collector service not linked.")
return
msg = await event.reply("⏳ Contacting Telegram to request login code...")
try:
sent = await self.collector.client.send_code_request(self.collector.phone)
self.collector.phone_code_hash = sent.phone_code_hash
await msg.edit(
f"📩 <b>Code Sent!</b>\n\n"
f"Telegram sent a verification code to <code>{self.collector.phone}</code>.\n\n"
f"Please reply with:\n"
f"<code>/code &lt;your_code&gt;</code>\n\n"
f"<i>Example:</i> <code>/code 12345</code>",
parse_mode="html"
)
except Exception as e:
await msg.edit(f"❌ Could not request code: {e}")
def _format_raw_post_caption(self, post: Post) -> str:
published_lines = ""
if post.published_to:
published_lines = "📤 <b>ارسال شده به کانال‌های:</b>\n"
for item in post.published_to:
if isinstance(item, dict):
t_title = item.get("target_title", "کانال مقصد")
published_lines += f" • <b>{t_title}</b>\n"
published_lines += "➖➖➖➖➖➖➖➖➖➖\n\n"
@self.client.on(events.NewMessage(pattern=r"^/code\s+(\S+)"))
async def cmd_code(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
await event.reply("Collector service not linked.")
return
code = event.pattern_match.group(1).strip()
status_msg = await event.reply("⏳ Submitting code to Telegram...")
result = await self.collector.submit_code(code)
await status_msg.edit(result, parse_mode="html", buttons=get_main_menu_keyboard())
caption = (
f"📥 <b>پست جدید از مبدا (<code>{post.source_channel_id}</code>):</b>\n\n"
f"{published_lines}"
f"{post.raw_text or ''}\n\n"
f"<i>👇 کانال مقصد مورد نظر را برای بازنویسی هوشمند انتخاب کنید:</i>"
)
return caption
@self.client.on(events.NewMessage(pattern=r"^/password\s+(.+)"))
async def cmd_password(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
await event.reply("Collector service not linked.")
return
pwd = event.pattern_match.group(1).strip()
status_msg = await event.reply("⏳ Verifying 2FA password...")
result = await self.collector.submit_password(pwd)
await status_msg.edit(result, parse_mode="html", buttons=get_main_menu_keyboard())
# --- Direct History Scraper Command ---
@self.client.on(events.NewMessage(pattern=r"^/scrape_history\s+(-?\d+)(?:\s+(\d+))?"))
async def cmd_scrape_history(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
await event.reply("Collector service not linked.")
return
ch_id = int(event.pattern_match.group(1))
limit = int(event.pattern_match.group(2)) if event.pattern_match.group(2) else 20
status_msg = await event.reply(f"⏳ Scraping the last <b>{limit}</b> posts from <code>{ch_id}</code> in the background...", parse_mode="html")
async def progress_notify(txt: str):
await status_msg.edit(txt, parse_mode="html", buttons=get_main_menu_keyboard())
await self.collector.scrape_channel_history(channel_id=ch_id, limit=limit, progress_callback=progress_notify)
# --- Statistics ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/stats|📊 Fleet Statistics)$"))
async def cmd_stats(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
pending_ai = len(await self.repo.get_posts_by_status("pending_ai", limit=1000))
pending_review = len(await self.repo.get_posts_by_status("pending_review", limit=1000))
approved = len(await self.repo.get_posts_by_status("approved", limit=1000))
published = len(await self.repo.get_posts_by_status("published", limit=1000))
rejected = len(await self.repo.get_posts_by_status("rejected", limit=1000))
redis_q = await self.collector.queue.qsize() if (self.collector and self.collector.queue) else 0
text = (
"📊 <b>Copykar Fleet Metrics</b>\n\n"
f"• 📥 <b>Redis Incoming Queue:</b> {redis_q} (Pacing: 1 post / 2m)\n"
f"• ⏳ <b>Pending AI:</b> {pending_ai}\n"
f"• 📋 <b>Pending Review:</b> {pending_review}\n"
f"• 🚀 <b>Approved (In Queue):</b> {approved}\n"
f"• ✅ <b>Published:</b> {published}\n"
f"• ❌ <b>Rejected:</b> {rejected}\n\n"
"📈 <i>Grafana Dashboard:</i> http://localhost:3000"
)
await event.reply(text, parse_mode="html", buttons=get_main_menu_keyboard())
# --- Sources Management with Interactive Scrape Buttons ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/sources|📡 Monitored Sources)$"))
async def cmd_sources(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
sources = await self.repo.get_active_sources()
if not sources:
await event.reply("No active source channels configured.\nTap <b> Add Source Guide</b> below to add one.", parse_mode="html", buttons=get_main_menu_keyboard())
return
await event.reply(f"📡 <b>Monitored Sources ({len(sources)} Active):</b>\nTap any button below to scrape past posts:", parse_mode="html")
for s in sources:
card = (
f"📢 <b>{s.title or 'Channel'}</b>\n"
f"• ID: <code>{s.channel_id}</code>\n"
f"• Username: @{s.username or 'none'}"
)
buttons = [
[
Button.inline(f"📥 Scrape 20 Posts", data=f"hist:{s.channel_id}:20"),
Button.inline(f"📥 Scrape 50 Posts", data=f"hist:{s.channel_id}:50"),
]
]
await event.reply(card, parse_mode="html", buttons=buttons)
@self.client.on(events.NewMessage(pattern=r"(?i)^( Add Source Guide)$"))
async def cmd_add_source_guide(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
guide = (
" <b>How to Add a Source Channel:</b>\n\n"
"Send the command in this format:\n"
"<code>/add_source &lt;channel_id&gt; &lt;title&gt; [username]</code>\n\n"
"<i>Example:</i>\n"
"<code>/add_source -1001234567890 TechNews technews_chan</code>"
)
await event.reply(guide, parse_mode="html")
@self.client.on(events.NewMessage(pattern=r"^/add_source\s+(-?\d+)\s+([^\s]+)(?:\s+([^\s]+))?"))
async def cmd_add_source(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
ch_id = int(event.pattern_match.group(1))
title = event.pattern_match.group(2)
username = event.pattern_match.group(3)
await self.repo.add_source(channel_id=ch_id, title=title, username=username)
buttons = [
[
Button.inline(f"📥 Scrape 20 Posts Now", data=f"hist:{ch_id}:20"),
Button.inline(f"📥 Scrape 50 Posts Now", data=f"hist:{ch_id}:50")
]
]
await event.reply(
f"✅ Added source channel <b>{title}</b> (<code>{ch_id}</code>).\n\nWould you like to scrape past posts now?",
parse_mode="html",
buttons=buttons
)
# --- Targets Management ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/targets|🎯 Target Channels)$"))
async def cmd_targets(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("No target channels configured.\nUse <code> Add Target Guide</code> to add one.", parse_mode="html")
return
lines = ["<b>🎯 Target Publishing Channels:</b>\n"]
for t in targets:
lines.append(f"• ID: <code>{t.id}</code> (Channel: <code>{t.channel_id}</code>)\n Title: <b>{t.title}</b> | Interval: <b>{t.post_interval_min}m</b>")
await event.reply("\n".join(lines), parse_mode="html", buttons=get_main_menu_keyboard())
@self.client.on(events.NewMessage(pattern=r"(?i)^( Add Target Guide)$"))
async def cmd_add_target_guide(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
guide = (
"🎯 <b>How to Add a Target Channel:</b>\n\n"
"Send the command in this format:\n"
"<code>/add_target &lt;channel_id&gt; &lt;title&gt; &lt;interval_minutes&gt; [username]</code>\n\n"
"<i>Example (posts every 30 minutes):</i>\n"
"<code>/add_target -1009876543210 MyMainChannel 30 my_main_chan</code>"
)
await event.reply(guide, parse_mode="html")
@self.client.on(events.NewMessage(pattern=r"^/add_target\s+(-?\d+)\s+([^\s]+)\s+(\d+)(?:\s+([^\s]+))?"))
async def cmd_add_target(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
ch_id = int(event.pattern_match.group(1))
title = event.pattern_match.group(2)
interval_min = int(event.pattern_match.group(3))
username = event.pattern_match.group(4)
await self.repo.add_target(channel_id=ch_id, title=title, username=username, post_interval_min=interval_min)
await event.reply(f"✅ Added target channel <b>{title}</b> with interval <b>{interval_min}m</b>.", parse_mode="html", buttons=get_main_menu_keyboard())
@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))
new_interval = int(event.pattern_match.group(2))
target = await self.repo.get_target_by_id(target_id)
if not target:
await event.reply("Target channel not found.")
return
await self.repo.add_target(
channel_id=target.channel_id,
title=target.title,
username=target.username,
post_interval_min=new_interval
)
await event.reply(f"✅ Updated interval for <b>{target.title}</b> to <b>{new_interval} minutes</b>.", parse_mode="html", buttons=get_main_menu_keyboard())
@self.client.on(events.NewMessage(pattern=r"(?i)^(/help|❓ Help & Documentation)$"))
async def cmd_help(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
help_text = (
"📖 <b>Copykar Bot Quick Help</b>\n\n"
"1. <b>Monitored Sources:</b> Tap <code>📡 Monitored Sources</code> to view channels and click <code>[📥 Scrape Posts]</code> on any channel.\n"
"2. <b>Review Flow:</b> AI scans posts, checks duplicates, and sends drafts to the review channel with inline approval buttons.\n"
"3. <b>Publishing:</b> Approved posts are published to your target channels strictly according to their interval minutes."
)
await event.reply(help_text, parse_mode="html", buttons=get_main_menu_keyboard())
# --- Inline Callback Queries ---
@self.client.on(events.CallbackQuery)
async def on_callback(event: events.CallbackQuery.Event):
if not self.is_admin(event.sender_id):
await event.answer("⛔ You are not authorized.", alert=True)
return
data = event.data.decode("utf-8")
# 1. Historical Scraping Callbacks
if data.startswith("hist:"):
_, ch_id_str, limit_str = data.split(":")
ch_id = int(ch_id_str)
limit = int(limit_str)
if not self.collector:
await event.answer("Collector service not linked.", alert=True)
return
await event.edit(f"⏳ <b>Scraping the last {limit} posts</b> from <code>{ch_id}</code>...", parse_mode="html", buttons=None)
async def progress_notify(txt: str):
await event.edit(txt, parse_mode="html")
await self.collector.scrape_channel_history(channel_id=ch_id, limit=limit, progress_callback=progress_notify)
await event.answer(f"Started scraping {limit} posts!")
# 2. Approval Callbacks
elif data.startswith("appr:"):
_, post_id_str, target_id_str = data.split(":")
post_id = int(post_id_str)
target_id = int(target_id_str)
target = await self.repo.get_target_by_id(target_id)
target_title = target.title if target else f"Target #{target_id}"
await self.repo.approve_post(post_id, target_id)
ADMIN_ACTIONS_TOTAL.labels(action="approved").inc()
await event.edit(
f"{event.text}\n\n✅ <b>Approved for {target_title}</b> by admin.",
parse_mode="html",
buttons=None
)
await event.answer(f"Approved for {target_title}!")
# 3. Reject Callbacks
elif data.startswith("rej:"):
_, post_id_str = data.split(":")
post_id = int(post_id_str)
await self.repo.reject_post(post_id)
ADMIN_ACTIONS_TOTAL.labels(action="rejected").inc()
await event.edit(
f"{event.text}\n\n❌ <b>Rejected</b> by admin.",
parse_mode="html",
buttons=None
)
await event.answer("Post rejected.")
async def send_review_post(self, post_id: int):
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:
return
targets = await self.repo.get_active_targets()
keyboard = get_review_keyboard(post.id, targets)
tags_str = ", ".join(post.tags) if post.tags else "None"
dup_warning = ""
if post.is_duplicate:
dup_warning = (
f"⚠️ <b>[DUPLICATE DETECTED]</b>\n"
f"<b>Reason:</b> {post.similarity_reason or 'Similar story already published'}\n"
f"<b>Matched Post ID:</b> #{post.duplicate_of_id}\n\n"
)
caption = (
f"📌 <b>Subject:</b> {post.subject or 'N/A'}\n"
f"🏷 <b>Tags:</b> <code>{tags_str}</code>\n\n"
f"{dup_warning}"
f"📝 <b>Generated Post Draft:</b>\n"
f"{post.ai_text or post.raw_text}\n\n"
f"<i>Source: Channel <code>{post.source_channel_id}</code> | Msg #{post.source_message_id}</i>"
)
keyboard = self._build_raw_post_keyboard(post, targets)
caption = self._format_raw_post_caption(post)
try:
if post.media_path and os.path.exists(post.media_path):
@@ -410,12 +142,441 @@ class AdminBotService:
parse_mode="html",
buttons=keyboard
)
await self.repo.update_review_message_id(post.id, msg.id)
except Exception as e:
logger.error(f"Failed to send review post {post.id} to review channel: {e}", exc_info=True)
logger.error(f"Failed to send raw post {post.id} to review channel: {e}", exc_info=True)
def _register_handlers(self):
# --- Start / Menu ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/start|/menu|منو)$"))
async def cmd_start(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
await event.reply(f"⛔ دسترسی غیرمجاز. شناسه عددی شما: <code>{event.sender_id}</code>", parse_mode="html")
return
userbot_status = "🔴 قطع / نیاز به لاگین"
if self.collector and self.collector.client.is_connected() and await self.collector.client.is_user_authorized():
me = await self.collector.client.get_me()
userbot_status = f"🟢 آنلاین ({me.first_name})"
welcome_text = (
"👋 <b>به پنل مدیریت سیستم هوشمند کپی‌کار خوش آمدید!</b>\n\n"
f"• 🤖 <b>وضعیت ربات جمع‌آوری‌کننده:</b> {userbot_status}\n"
f"• 📋 <b>شناسه کانال ادمین‌ها:</b> <code>{self.review_channel_id}</code>\n\n"
"از دکمه‌های زیر برای مدیریت کانال‌ها، تنظیم شخصیت و آمار استفاده کنید:"
)
await event.reply(welcome_text, parse_mode="html", buttons=get_persian_main_menu())
# --- Statistics ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/stats|📊 آمار و وضعیت ناوگان)$"))
async def cmd_stats(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
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
text = (
"📊 <b>آمار زنده سیستم کپی‌کار:</b>\n\n"
f"• 📥 <b>پست‌های موجود در صف ردیس:</b> <b>{redis_q}</b>\n"
f"• 📋 <b>پست‌های در انتظار بررسی ادمین:</b> <b>{pending_review}</b>\n"
f"• 🚀 <b>پست‌های منتشر شده:</b> <b>{published}</b>\n"
f"• ❌ <b>پست‌های رد شده:</b> <b>{rejected}</b>\n\n"
"📈 <i>داشبورد مانیتورینگ گرانافا:</i> http://localhost:3000"
)
await event.reply(text, parse_mode="html", buttons=get_persian_main_menu())
# --- Userbot Authentication ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/request_code|🔑 درخواست کد لاگین)$"))
async def cmd_request_code(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
await event.reply("سرویس کالکتور متصل نیست.")
return
msg = await event.reply("⏳ در حال ارسال درخواست کد لاگین به تلگرام...")
try:
sent = await self.collector.client.send_code_request(self.collector.phone)
self.collector.phone_code_hash = sent.phone_code_hash
await msg.edit(
f"📩 <b>کد تایید ارسال شد!</b>\n\n"
f"کد ارسال شده به شماره <code>{self.collector.phone}</code> را به این صورت ارسال کنید:\n"
f"<code>/code 12345</code>",
parse_mode="html"
)
except Exception as e:
await msg.edit(f"❌ خطا در درخواست کد: {e}")
@self.client.on(events.NewMessage(pattern=r"^/code\s+(\S+)"))
async def cmd_code(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
return
code = event.pattern_match.group(1).strip()
msg = await event.reply("⏳ در حال بررسی کد تایید...")
result = await self.collector.submit_code(code)
await msg.edit(result, parse_mode="html", buttons=get_persian_main_menu())
@self.client.on(events.NewMessage(pattern=r"^/password\s+(.+)"))
async def cmd_password(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
if not self.collector:
return
pwd = event.pattern_match.group(1).strip()
msg = await event.reply("⏳ در حال تایید رمز دو مرحله‌ای...")
result = await self.collector.submit_password(pwd)
await msg.edit(result, parse_mode="html", buttons=get_persian_main_menu())
# --- Sources Management ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/sources|📡 کانال‌های مبدا)$"))
async def cmd_sources(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
sources = await self.repo.get_active_sources()
if not sources:
await event.reply("هیچ کانال مبدایی ثبت نشده است. از دکمه <b>➕ افزودن کانال مبدا</b> استفاده کنید.", parse_mode="html")
return
await event.reply(f"📡 <b>کانال‌های مبدا فعال ({len(sources)} کانال):</b>", parse_mode="html")
for s in sources:
card = (
f"📢 <b>{s.title or 'کانال'}</b>\n"
f"• شناسه: <code>{s.channel_id}</code>\n"
f"• یوزرنیم: @{s.username or 'ندارد'}"
)
buttons = [
[
Button.inline("📥 استخراج ۲۰ پست گذشته", data=f"hist:{s.channel_id}:20"),
Button.inline("📥 استخراج ۵۰ پست گذشته", data=f"hist:{s.channel_id}:50"),
]
]
await event.reply(card, parse_mode="html", buttons=buttons)
@self.client.on(events.NewMessage(pattern=r"(?i)^(➕ افزودن کانال مبدا)$"))
async def cmd_add_source_guide(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
guide = (
"➕ <b>راهنمای افزودن کانال مبدا:</b>\n\n"
"دستور را با فرمت زیر ارسال کنید:\n"
"<code>/add_source &lt;شناسه_عددی_کانال&gt; &lt;عنوان&gt; [یوزرنیم]</code>\n\n"
"<i>مثال:</i>\n"
"<code>/add_source -1001234567890 اخبار_فوری fouri_news</code>"
)
await event.reply(guide, parse_mode="html")
@self.client.on(events.NewMessage(pattern=r"^/add_source\s+(-?\d+)\s+([^\s]+)(?:\s+([^\s]+))?"))
async def cmd_add_source(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
ch_id = int(event.pattern_match.group(1))
title = event.pattern_match.group(2)
username = event.pattern_match.group(3)
await self.repo.add_source(channel_id=ch_id, title=title, username=username)
buttons = [
[
Button.inline("📥 استخراج ۲۰ پست گذشته این کانال", data=f"hist:{ch_id}:20"),
Button.inline("📥 استخراج ۵۰ پست گذشته این کانال", data=f"hist:{ch_id}:50")
]
]
await event.reply(
f"✅ کانال مبدا <b>{title}</b> (<code>{ch_id}</code>) با موفقیت افزوده شد.\nآیا می‌خواهید پست‌های قبلی این کانال را هم دریافت کنید؟",
parse_mode="html",
buttons=buttons
)
# --- Targets Management ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/targets|🎯 کانال‌های مقصد)$"))
async def cmd_targets(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("هیچ کانال مقصدی ثبت نشده است. از دکمه <b>➕ افزودن کانال مقصد</b> استفاده کنید.", parse_mode="html")
return
lines = ["<b>🎯 کانال‌های مقصد برای انتشار:</b>\n"]
for t in targets:
lines.append(
f"• <b>{t.title}</b> (شناسه: <code>{t.channel_id}</code> | ID دیتابیس: <code>{t.id}</code>)\n"
f" 🎭 <b>شخصیت و لحن:</b> {t.personality or 'پیش‌فرض'}\n"
f" 🏷 <b>فوتر / تگ‌ها:</b> <code>{t.custom_footer or 'ندارد'}</code>\n"
)
await event.reply("\n".join(lines), parse_mode="html", buttons=get_persian_main_menu())
@self.client.on(events.NewMessage(pattern=r"(?i)^(➕ افزودن کانال مقصد)$"))
async def cmd_add_target_guide(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
guide = (
"🎯 <b>راهنمای افزودن کانال مقصد:</b>\n\n"
"دستور را با فرمت زیر ارسال کنید:\n"
"<code>/add_target &lt;شناسه_کانال&gt; &lt;عنوان&gt; [یوزرنیم]</code>\n\n"
"<i>مثال:</i>\n"
"<code>/add_target -1009876543210 دنیای_هوش_مصنوعی ai_world_chan</code>"
)
await event.reply(guide, parse_mode="html")
@self.client.on(events.NewMessage(pattern=r"^/add_target\s+(-?\d+)\s+([^\s]+)(?:\s+([^\s]+))?"))
async def cmd_add_target(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
ch_id = int(event.pattern_match.group(1))
title = event.pattern_match.group(2)
username = event.pattern_match.group(3)
tid = await self.repo.add_target(channel_id=ch_id, title=title, username=username)
await event.reply(
f"✅ کانال مقصد <b>{title}</b> افزوده شد (ID دیتابیس: <code>{tid}</code>).\n\n"
f"اکنون می‌توانید با دکمه <b>🎭 تنظیم شخصیت کانال‌ها</b> لحن و تگ‌های آن را تنظیم کنید.",
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):
if not self.is_admin(event.sender_id):
return
targets = await self.repo.get_active_targets()
if not targets:
await event.reply("ابتدا با استفاده از <b>➕ افزودن کانال مقصد</b> یک کانال مقصد اضافه کنید.", parse_mode="html")
return
text = (
"🎭 <b>تنظیم شخصیت، لحن و تگ‌های کانال‌های مقصد:</b>\n\n"
"برای تغییر لحن و استایل نگارش کانال از دستور زیر استفاده کنید:\n"
"<code>/set_personality &lt;شناسه_دیتابیس_کانال&gt; &lt;توضیحات لحن&gt;</code>\n\n"
"<i>مثال:</i>\n"
"<code>/set_personality 1 لحن جذاب و ژورنالیستی، استفاده از تیترهای بولد و ایموجی‌های مرتبط</code>\n\n"
"برای تنظیم فوتر و هشتگ‌های اختصاصی انتهای پست:\n"
"<code>/set_footer &lt;شناسه_دیتابیس_کانال&gt; &lt;تگ‌ها یا آیدی کانال&gt;</code>\n\n"
"<i>مثال:</i>\n"
"<code>/set_footer 1 🆔 @my_tech_chan\n#تکنولوژی #هوش_مصنوعی</code>\n\n"
"<b>کانال‌های موجود و شخصیت فعلی:</b>\n"
)
for t in targets:
text += (
f"• ID: <code>{t.id}</code> | <b>{t.title}</b>\n"
f" 🎭 لحن: <i>{t.personality or 'پیش‌فرض'}</i>\n"
f" 🏷 فوتر: <code>{t.custom_footer or 'ندارد'}</code>\n\n"
)
await event.reply(text, parse_mode="html", buttons=get_persian_main_menu())
@self.client.on(events.NewMessage(pattern=r"^/set_personality\s+(\d+)\s+(.+)"))
async def cmd_set_personality(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
target_id = int(event.pattern_match.group(1))
personality = event.pattern_match.group(2).strip()
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_personality(target_id, personality)
await event.reply(
f"✅ شخصیت و لحن کانال <b>{target.title}</b> با موفقیت به روز شد:\n\n<i>{personality}</i>",
parse_mode="html",
buttons=get_persian_main_menu()
)
@self.client.on(events.NewMessage(pattern=r"^/set_footer\s+(\d+)\s+([\s\S]+)"))
async def cmd_set_footer(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
target_id = int(event.pattern_match.group(1))
footer = event.pattern_match.group(2).strip()
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_footer(target_id, footer)
await event.reply(
f"✅ فوتر اختصاصی کانال <b>{target.title}</b> به روز شد:\n\n<code>{footer}</code>",
parse_mode="html",
buttons=get_persian_main_menu()
)
# --- Help ---
@self.client.on(events.NewMessage(pattern=r"(?i)^(/help|❓ راهنمای سیستم)$"))
async def cmd_help(event: events.NewMessage.Event):
if not self.is_admin(event.sender_id):
return
help_text = (
"📖 <b>راهنمای فرآیند کاری سیستم کپی‌کار:</b>\n\n"
"1. 📥 <b>دریافت خام پست‌ها:</b> پست‌ها بدون پردازش هوش مصنوعی مستقیماً به کانال ادمین‌ها ارسال می‌شوند.\n"
"2. 🎯 <b>انتخاب کانال مقصد:</b> با لمس دکمه هر کانال، هوش مصنوعی پست را متناسب با شخصیت، استایل و فوتر اختصاصی همان کانال بازنویسی کرده و تمام تگ‌ها و لینک‌های مبدا را حذف می‌کند.\n"
"3. 👁 <b>پیش‌نمایش زنده:</b> پیش‌نمایش بازنویسی شده به همراه دکمه تایید نهایی نمایش داده می‌شود.\n"
"4. 🚀 <b>انتشار و ارسال مجدد:</b> پس از انتشار، پست اصلی در کانال ادمین بازگردانده شده و سابقه انتشار نمایش می‌یابد تا بتوانید آن را به سایر کانال‌ها نیز ارسال کنید."
)
await event.reply(help_text, parse_mode="html", buttons=get_persian_main_menu())
# --- Interactive Inline Callbacks for Reviews & Target Rewrites ---
@self.client.on(events.CallbackQuery)
async def on_callback(event: events.CallbackQuery.Event):
if not self.is_admin(event.sender_id):
await event.answer("⛔ دسترسی غیرمجاز.", alert=True)
return
data = event.data.decode("utf-8")
# 1. Historical Scraping Callback
if data.startswith("hist:"):
_, ch_id_str, limit_str = data.split(":")
ch_id = int(ch_id_str)
limit = int(limit_str)
if not self.collector:
await event.answer("کالکتور در دسترس نیست.", alert=True)
return
await event.edit(f"⏳ <b>در حال دریافت {limit} پست گذشته</b> از کانال <code>{ch_id}</code>...", parse_mode="html")
async def progress_notify(txt: str):
await event.edit(txt, parse_mode="html")
await self.collector.scrape_channel_history(channel_id=ch_id, limit=limit, progress_callback=progress_notify)
await event.answer("فرآیند دریافت آغاز شد.")
# 2. Target Selected -> Trigger On-Demand AI Rewrite for that Target
elif data.startswith("sel_trg:"):
_, post_id_str, target_id_str = data.split(":")
post_id = int(post_id_str)
target_id = int(target_id_str)
post = await self.repo.get_post_by_id(post_id)
target = await self.repo.get_target_by_id(target_id)
if not post or not target:
await event.answer("پست یا کانال مقصد یافت نشد.", alert=True)
return
await event.answer(f"در حال بازنویسی برای {target.title}...")
# Show loading placeholder
loading_caption = (
f"🤖 <b>در حال بازنویسی هوشمند برای کانال:</b> <b>{target.title}</b>...\n"
f"<i>(اعمال لحن اختصاصی و حذف تگ‌های مبدا)</i>"
)
try:
await event.edit(loading_caption, parse_mode="html", buttons=None)
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"🎯 <b>پیش‌نمایش بازنویسی شده برای: {target.title}</b>\n"
f"🎭 <b>شخصیت و لحن:</b> <i>{target.personality or 'پیش‌فرض'}</i>\n"
f"➖➖➖➖➖➖➖➖➖➖\n\n"
f"{rewritten_text}\n\n"
f"➖➖➖➖➖➖➖➖➖➖\n"
f"<i>آیا این متن مورد تایید است؟</i>"
)
preview_buttons = [
[
Button.inline(f"✅ تایید و ارسال به {target.title}", data=f"pub:{post_id}:{target_id}"),
],
[
Button.inline("🔙 انصراف / بازگشت به پست اصلی", data=f"cancel:{post_id}")
]
]
await event.edit(preview_caption, parse_mode="html", buttons=preview_buttons)
# 3. Publish to Target Confirmed
elif data.startswith("pub:"):
_, post_id_str, target_id_str = data.split(":")
post_id = int(post_id_str)
target_id = int(target_id_str)
post = await self.repo.get_post_by_id(post_id)
target = await self.repo.get_target_by_id(target_id)
if not post or not target:
await event.answer("اطلاعات یافت نشد.", alert=True)
return
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}...")
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"
)
# 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)
await event.edit(
f"✅ <b>با موفقیت در {target.title} منتشر شد!</b>\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)
# 4. Cancel Preview & Restore Original Card
elif data.startswith("cancel:"):
_, post_id_str = data.split(":")
post_id = int(post_id_str)
post = await self.repo.get_post_by_id(post_id)
if not post:
return
targets = await self.repo.get_active_targets()
caption = self._format_raw_post_caption(post)
buttons = self._build_raw_post_keyboard(post, targets)
await event.edit(caption, parse_mode="html", buttons=buttons)
await event.answer("پیش‌نمایش لغو شد.")
# 5. Reject Post
elif data.startswith("rej:"):
_, post_id_str = data.split(":")
post_id = int(post_id_str)
await self.repo.reject_post(post_id)
ADMIN_ACTIONS_TOTAL.labels(action="rejected").inc()
await event.edit(
f"{event.text}\n\n❌ <b>این پست توسط ادمین رد و بایگانی شد.</b>",
parse_mode="html",
buttons=None
)
await event.answer("پست بایگانی شد.")
async def stop(self):
if self.client.is_connected():
await self.client.disconnect()
logger.info("Admin Review Bot disconnected.")
logger.info("Admin Bot disconnected.")
+44 -116
View File
@@ -1,47 +1,30 @@
import logging
import json
from typing import List, Optional, Dict, Any
from db.models import Post, TargetChannel
from db.repository import Repository
from core.llm import LLMClient
from core.dedup import compute_content_hash
from core.metrics import DUPLICATES_DETECTED_TOTAL
logger = logging.getLogger(__name__)
TAG_EXTRACTION_SYSTEM_PROMPT = """
You are an AI news analyst and classifier.
Given a social media/channel post, extract:
1. "subject": A brief, specific headline/subject (3-8 words).
2. "tags": A JSON array of 3 to 6 lowercase keywords/topics/entities (e.g. ["ai", "nvidia", "gpus", "hardware"]).
Respond ONLY in JSON format:
CHANNEL_REWRITE_SYSTEM_PROMPT = """
You are a professional Persian Telegram copywriter and editor.
Your job is to rewrite the provided raw post specifically for the target channel: "{channel_title}".
CHANNEL PERSONALITY & TONE GUIDELINES:
{personality}
CRITICAL RULES:
1. Completely REMOVE all original channel usernames (e.g. @source_channel), sponsor tags, author watermarks, and source links.
2. Translate or rewrite into natural, highly engaging, and fluent Persian (فارسی روان، جذاب و حرفه‌ای).
3. Use appropriate emojis and clear paragraph spacing.
4. If a custom footer/tag is provided below, append it cleanly at the very end of the post:
{custom_footer}
Respond ONLY in valid JSON format:
{
"subject": "...",
"tags": ["tag1", "tag2", "tag3"]
}
"""
DUPLICATE_CHECK_SYSTEM_PROMPT = """
You are an expert news editor checking for duplicate news stories.
Given a NEW POST and a list of PREVIOUS POSTS, determine if the NEW POST is covering the same exact event, news item, or story as any of the previous posts.
Respond ONLY in JSON format:
{
"is_duplicate": true/false,
"duplicate_of_id": <id of matched previous post or null>,
"similarity_reason": "<short explanation of why it is or is not a duplicate>"
}
"""
POST_REWRITE_SYSTEM_PROMPT = """
You are an expert Telegram content creator and copywriter.
Rewrite the provided post to make it engaging, well-formatted, professional, and clear.
Use appropriate emojis, clear paragraphs, and markdown formatting.
Remove any original promotional links, author credits, or watermarks.
Respond ONLY in JSON format:
{
"ai_text": "...",
"suggested_target_id": <optional id of best matching target channel or null>
"rewritten_text": "..."
}
"""
@@ -50,89 +33,34 @@ class AIProcessor:
self.repo = repo
self.llm = llm or LLMClient()
async def process_post(self, post_id: int) -> Optional[Post]:
post = await self.repo.get_post_by_id(post_id)
if not post or not post.raw_text:
return post
async def rewrite_for_target(self, raw_text: str, target: TargetChannel) -> str:
"""Rewrite raw text according to a specific target channel's personality and custom footer."""
if not raw_text:
return ""
raw_text = post.raw_text
is_dup = False
dup_of_id = None
sim_reason = None
personality_text = target.personality.strip() if target.personality else "لحن رسمی، جذاب و روان به همراه ایموجی‌های مرتبط و پاراگراف‌بندی مرتب."
footer_text = target.custom_footer.strip() if target.custom_footer else (f"@{target.username}" if target.username else "")
# 1. Exact hash duplicate check
content_hash = post.content_hash or compute_content_hash(raw_text)
if content_hash:
exact_dup = await self.repo.find_duplicate_post_by_hash(content_hash)
if exact_dup and exact_dup.id != post.id:
is_dup = True
dup_of_id = exact_dup.id
sim_reason = "Exact match on normalized text/media hash"
DUPLICATES_DETECTED_TOTAL.labels(method="hash").inc()
# 2. Extract Tags and Subject via AI
tags = []
subject = "General News"
try:
tag_res = await self.llm.generate_json(
prompt=f"Post content:\n\n{raw_text}",
system_prompt=TAG_EXTRACTION_SYSTEM_PROMPT,
action_name="extract_tags"
)
subject = tag_res.get("subject", subject)
tags = [t.lower().strip() for t in tag_res.get("tags", []) if isinstance(t, str)]
await self.repo.update_post_tags(post.id, tags, subject)
except Exception as e:
logger.error(f"Tag extraction failed for post {post.id}: {e}")
# 3. Candidate search & Semantic AI Deduplication check (if not already exact dup)
if not is_dup and tags:
candidates = await self.repo.find_candidate_posts_by_tags(tags, exclude_post_id=post.id, hours_lookback=72, limit=5)
if candidates:
cand_texts = "\n---\n".join([f"ID {c.id} (Subject: {c.subject}):\n{c.raw_text}" for c in candidates if c.raw_text])
prompt = f"NEW POST:\n{raw_text}\n\nPREVIOUS CANDIDATE POSTS:\n{cand_texts}"
try:
dup_res = await self.llm.generate_json(
prompt=prompt,
system_prompt=DUPLICATE_CHECK_SYSTEM_PROMPT,
action_name="check_duplicate"
)
if dup_res.get("is_duplicate"):
is_dup = True
dup_of_id = dup_res.get("duplicate_of_id")
sim_reason = dup_res.get("similarity_reason", "AI detected duplicate news topic")
DUPLICATES_DETECTED_TOTAL.labels(method="ai_semantic").inc()
except Exception as e:
logger.error(f"Semantic duplicate check failed for post {post.id}: {e}")
# 4. Rewrite post for our channels
ai_text = raw_text
suggested_target_id = None
targets = await self.repo.get_active_targets()
target_info = "\n".join([f"Target ID {t.id}: {t.title} (@{t.username or 'none'})" for t in targets])
rewrite_prompt = f"TARGET CHANNELS AVAILABLE:\n{target_info or 'None'}\n\nORIGINAL POST:\n{raw_text}"
try:
rewrite_res = await self.llm.generate_json(
prompt=rewrite_prompt,
system_prompt=POST_REWRITE_SYSTEM_PROMPT,
action_name="rewrite_post"
)
ai_text = rewrite_res.get("ai_text", raw_text)
suggested_target_id = rewrite_res.get("suggested_target_id")
except Exception as e:
logger.error(f"Post rewrite failed for post {post.id}: {e}")
# 5. Save AI results into database
await self.repo.update_ai_result(
post_id=post.id,
subject=subject,
ai_text=ai_text,
tags=tags,
suggested_target_id=suggested_target_id,
is_duplicate=is_dup,
duplicate_of_id=dup_of_id,
similarity_reason=sim_reason,
sys_prompt = CHANNEL_REWRITE_SYSTEM_PROMPT.format(
channel_title=target.title or "کانال تلگرام",
personality=personality_text,
custom_footer=footer_text
)
return await self.repo.get_post_by_id(post.id)
try:
res = await self.llm.generate_json(
prompt=f"متن اصلی پست برای بازنویسی:\n\n{raw_text}",
system_prompt=sys_prompt,
action_name="rewrite_target_post"
)
rewritten = res.get("rewritten_text")
if rewritten:
return rewritten.strip()
except Exception as e:
logger.error(f"Failed to rewrite post for target {target.id} ({target.title}): {e}")
# Fallback if AI fails: clean basic @mentions and append footer
fallback = raw_text
if footer_text:
fallback = f"{fallback}\n\n{footer_text}"
return fallback
+28 -23
View File
@@ -6,9 +6,8 @@ from telethon.errors import SessionPasswordNeededError
from telethon.tl.types import MessageMediaPhoto, MessageMediaDocument
from db.repository import Repository
from core.dedup import compute_content_hash, compute_file_hash
from services.ai_processor import AIProcessor
from core.queue import RedisQueue
from core.metrics import COLLECTED_POSTS_TOTAL
from core.metrics import COLLECTED_POSTS_TOTAL, SOURCE_ACTIVITY_TOTAL
from core.proxy import get_telegram_proxy
logger = logging.getLogger(__name__)
@@ -20,7 +19,7 @@ class CollectorService:
def __init__(
self,
repo: Repository,
ai_processor: AIProcessor,
on_post_received: Optional[Callable[[int], Awaitable[None]]] = None,
queue: Optional[RedisQueue] = None,
api_id: Optional[int] = None,
api_hash: Optional[str] = None,
@@ -28,7 +27,7 @@ class CollectorService:
session_name: Optional[str] = None,
):
self.repo = repo
self.ai_processor = ai_processor
self.on_post_received = on_post_received
self.queue = queue
self.api_id = api_id or int(os.getenv("API_ID", "0"))
self.api_hash = api_hash or os.getenv("API_HASH", "")
@@ -56,14 +55,14 @@ class CollectorService:
sent = await self.client.send_code_request(self.phone)
self.phone_code_hash = sent.phone_code_hash
await notify_fn(
f"🔐 <b>Collector Userbot Login Required</b>\n\n"
f"A login code was sent to phone <code>{self.phone}</code>.\n\n"
f"Please reply with: <code>/code &lt;your_code&gt;</code>\n"
f"(Or <code>/password &lt;2fa_password&gt;</code> if 2FA is enabled)."
f"🔐 <b>نیاز به ورود ربات جمع‌آوری‌کننده</b>\n\n"
f"کد تایید تلگرام به شماره <code>{self.phone}</code> ارسال شد.\n\n"
f"لطفا با دستور زیر پاسخ دهید:\n"
f"<code>/code 12345</code>"
)
except Exception as e:
logger.error(f"Failed to send login code request: {e}")
await notify_fn(f"Failed to request login code: {e}")
await notify_fn(f"خطا در ارسال کد ورود: {e}")
return False
async def submit_code(self, code: str) -> str:
@@ -75,20 +74,20 @@ class CollectorService:
await self.client.sign_in(phone=self.phone, code=code, phone_code_hash=self.phone_code_hash)
me = await self.client.get_me()
self._register_handlers()
return f"Logged in successfully as <b>{me.first_name}</b> (@{me.username or 'none'}). Collector is now active!"
return f"ورود موفقیت‌آمیز بود! حساب فعال: <b>{me.first_name}</b> (@{me.username or 'ندارد'})."
except SessionPasswordNeededError:
return "🔐 <b>Two-Factor Authentication (2FA) is enabled.</b> Please send: <code>/password &lt;your_2fa_password&gt;</code>"
return "🔐 <b>رمز دو مرحله‌ای فعال است.</b> لطفا با این دستور رمز را وارد کنید: <code>/password رمز_عبور</code>"
except Exception as e:
return f"Login failed: {e}"
return f"خطا در ورود: {e}"
async def submit_password(self, password: str) -> str:
try:
await self.client.sign_in(password=password)
me = await self.client.get_me()
self._register_handlers()
return f"2FA Verified! Logged in as <b>{me.first_name}</b> (@{me.username or 'none'}). Collector is now active!"
return f"تایید دو مرحله‌ای موفق بود! حساب فعال: <b>{me.first_name}</b>."
except Exception as e:
return f"2FA verification failed: {e}"
return f"خطا در تایید رمز دو مرحله‌ای: {e}"
def _register_handlers(self):
if self._handlers_registered:
@@ -141,11 +140,13 @@ class CollectorService:
if post_id:
COLLECTED_POSTS_TOTAL.labels(source_channel_id=str(chat_id)).inc()
logger.info(f"Collected new post ID {post_id} from channel {chat_id}")
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)
else:
await self.ai_processor.process_post(post_id)
elif 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)
@@ -158,9 +159,11 @@ class CollectorService:
"""Scrape historical messages from a source channel."""
if not self.client.is_connected() or not await self.client.is_user_authorized():
if progress_callback:
await progress_callback("Collector Userbot is not authorized. Please log in first.")
await progress_callback("ربات متصل نیست. لطفا ابتدا لاگین کنید.")
return 0
source = await self.repo.get_source_by_channel_id(channel_id)
source_title = source.title if source else str(channel_id)
collected_count = 0
skipped_count = 0
@@ -210,23 +213,25 @@ class CollectorService:
if post_id:
collected_count += 1
COLLECTED_POSTS_TOTAL.labels(source_channel_id=str(channel_id)).inc()
logger.info(f"Backfilled historical post ID {post_id} from {channel_id}")
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)
else:
await self.ai_processor.process_post(post_id)
elif self.on_post_received:
await self.on_post_received(post_id)
else:
skipped_count += 1
if progress_callback:
await progress_callback(
f"Scraped <b>{collected_count}</b> new posts from <code>{channel_id}</code> and queued in Redis! (Skipped {skipped_count} existing/empty)."
f"تعداد <b>{collected_count}</b> پست جدید از <code>{channel_id}</code> دریافت و در کانال ادمین قرار گرفت! (رد شده تکراری: {skipped_count})."
)
return collected_count
except Exception as e:
logger.error(f"Error scraping history from {channel_id}: {e}", exc_info=True)
if progress_callback:
await progress_callback(f"Error scraping channel <code>{channel_id}</code>: {e}")
await progress_callback(f"خطا در دریافت پست‌های کانال <code>{channel_id}</code>: {e}")
return collected_count
async def stop(self):
+5 -6
View File
@@ -3,7 +3,6 @@ import asyncio
import logging
from typing import Optional
from core.queue import RedisQueue
from services.ai_processor import AIProcessor
from core.metrics import QUEUE_POSTS_GAUGE, REDIS_QUEUE_SIZE_GAUGE
logger = logging.getLogger(__name__)
@@ -11,9 +10,9 @@ logger = logging.getLogger(__name__)
FETCH_INTERVAL_SECONDS = int(os.getenv("AI_PROCESSING_INTERVAL_SECONDS", "120"))
class QueueConsumerService:
def __init__(self, queue: RedisQueue, ai_processor: AIProcessor, fetch_interval: int = FETCH_INTERVAL_SECONDS):
def __init__(self, queue: RedisQueue, on_post_popped = None, fetch_interval: int = FETCH_INTERVAL_SECONDS):
self.queue = queue
self.ai_processor = ai_processor
self.on_post_popped = on_post_popped
self.fetch_interval = fetch_interval
self._running = False
self._task: Optional[asyncio.Task] = None
@@ -33,9 +32,9 @@ class QueueConsumerService:
if qsize > 0:
post_id = await self.queue.pop()
if post_id:
logger.info(f"Paced Consumer: processing post ID {post_id} from Redis queue (remaining: {qsize - 1})")
await self.ai_processor.process_post(post_id)
if post_id and self.on_post_popped:
logger.info(f"Dispatching raw post ID {post_id} to Admin Review channel (remaining in Redis: {qsize - 1})")
await self.on_post_popped(post_id)
new_qsize = await self.queue.qsize()
QUEUE_POSTS_GAUGE.labels(status="redis_incoming").set(new_qsize)
REDIS_QUEUE_SIZE_GAUGE.set(new_qsize)