relayers
This commit is contained in:
1 parent
21964fa986
commit
797f379648
68 files changed
+23871
-1271
No files matched your search
@@ -1,267 +1,602 @@
|
||||
"""Media conversion service for processing uploaded files."""
|
||||
|
||||
import asyncio
|
||||
from datetime import datetime
|
||||
import os
|
||||
import uuid
|
||||
import hashlib
|
||||
import json
|
||||
import shutil
|
||||
import magic # python-magic for MIME detection
|
||||
from base58 import b58decode, b58encode
|
||||
from sqlalchemy import and_, or_
|
||||
from app.core.models.node_storage import StoredContent
|
||||
from app.core.models._telegram import Wrapped_CBotChat
|
||||
from app.core._utils.send_status import send_status
|
||||
from app.core.logger import make_log
|
||||
from app.core.models.user import User
|
||||
from app.core.models import WalletConnection
|
||||
from app.core.storage import db_session
|
||||
from app.core._config import UPLOADS_DIR
|
||||
from app.core.content.content_id import ContentId
|
||||
import logging
|
||||
import os
|
||||
import tempfile
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional, Set, Any, Tuple
|
||||
|
||||
import aiofiles
|
||||
import redis.asyncio as redis
|
||||
from PIL import Image, ImageOps
|
||||
from sqlalchemy import select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.core.database import get_async_session
|
||||
from app.core.models.content import Content, FileUpload
|
||||
from app.core.storage import storage_manager
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def convert_loop(memory):
|
||||
with db_session() as session:
|
||||
# Query for unprocessed encrypted content
|
||||
unprocessed_encrypted_content = session.query(StoredContent).filter(
|
||||
and_(
|
||||
StoredContent.type == "onchain/content",
|
||||
or_(
|
||||
StoredContent.btfs_cid == None,
|
||||
StoredContent.ipfs_cid == None,
|
||||
)
|
||||
)
|
||||
).first()
|
||||
if not unprocessed_encrypted_content:
|
||||
make_log("ConvertProcess", "No content to convert", level="debug")
|
||||
return
|
||||
class ConvertService:
|
||||
"""Service for converting and processing uploaded media files."""
|
||||
|
||||
# Достаем расшифрованный файл
|
||||
decrypted_content = session.query(StoredContent).filter(
|
||||
StoredContent.id == unprocessed_encrypted_content.decrypted_content_id
|
||||
).first()
|
||||
if not decrypted_content:
|
||||
make_log("ConvertProcess", "Decrypted content not found", level="error")
|
||||
return
|
||||
|
||||
# Определяем путь и расширение входного файла
|
||||
input_file_path = f"/Storage/storedContent/{decrypted_content.hash}"
|
||||
input_ext = (unprocessed_encrypted_content.filename.split('.')[-1]
|
||||
if '.' in unprocessed_encrypted_content.filename else "mp4")
|
||||
|
||||
# ==== Новая логика: определение MIME-тип через python-magic ====
|
||||
def __init__(self):
|
||||
self.settings = get_settings()
|
||||
self.redis_client: Optional[redis.Redis] = None
|
||||
self.is_running = False
|
||||
self.tasks: Set[asyncio.Task] = set()
|
||||
|
||||
# Conversion configuration
|
||||
self.batch_size = 10
|
||||
self.process_interval = 5 # seconds
|
||||
self.max_retries = 3
|
||||
self.temp_dir = Path(tempfile.gettempdir()) / "uploader_convert"
|
||||
self.temp_dir.mkdir(exist_ok=True)
|
||||
|
||||
# Supported formats
|
||||
self.image_formats = {'.jpg', '.jpeg', '.png', '.gif', '.bmp', '.webp', '.tiff'}
|
||||
self.video_formats = {'.mp4', '.avi', '.mov', '.wmv', '.flv', '.webm', '.mkv'}
|
||||
self.audio_formats = {'.mp3', '.wav', '.ogg', '.m4a', '.flac', '.aac'}
|
||||
self.document_formats = {'.pdf', '.doc', '.docx', '.txt', '.rtf'}
|
||||
|
||||
# Image processing settings
|
||||
self.thumbnail_sizes = [(150, 150), (300, 300), (800, 600)]
|
||||
self.image_quality = 85
|
||||
self.max_image_size = (2048, 2048)
|
||||
|
||||
async def start(self) -> None:
|
||||
"""Start the conversion service."""
|
||||
try:
|
||||
mime_type = magic.from_file(input_file_path.replace("/Storage/storedContent", "/app/data"), mime=True)
|
||||
except Exception as e:
|
||||
make_log("ConvertProcess", f"magic probe failed: {e}", level="warning")
|
||||
mime_type = ""
|
||||
|
||||
if mime_type.startswith("video/"):
|
||||
content_kind = "video"
|
||||
elif mime_type.startswith("audio/"):
|
||||
content_kind = "audio"
|
||||
else:
|
||||
content_kind = "other"
|
||||
|
||||
make_log("ConvertProcess", f"Detected content_kind={content_kind}, mime={mime_type}", level="info")
|
||||
|
||||
# Для прочих типов сохраняем raw копию и выходим
|
||||
if content_kind == "other":
|
||||
make_log("ConvertProcess", f"Content {unprocessed_encrypted_content.id} processed. Not audio/video, copy just", level="info")
|
||||
unprocessed_encrypted_content.btfs_cid = ContentId(
|
||||
version=2, content_hash=b58decode(decrypted_content.hash)
|
||||
).serialize_v2()
|
||||
unprocessed_encrypted_content.ipfs_cid = ContentId(
|
||||
version=2, content_hash=b58decode(decrypted_content.hash)
|
||||
).serialize_v2()
|
||||
unprocessed_encrypted_content.meta = {
|
||||
**unprocessed_encrypted_content.meta,
|
||||
'converted_content': {
|
||||
option_name: decrypted_content.hash for option_name in ['high', 'low', 'low_preview']
|
||||
}
|
||||
}
|
||||
session.commit()
|
||||
return
|
||||
|
||||
# ==== Конвертация для видео или аудио: оригинальная логика ====
|
||||
# Static preview interval in seconds
|
||||
preview_interval = [0, 30]
|
||||
if unprocessed_encrypted_content.onchain_index in [2]:
|
||||
preview_interval = [0, 60]
|
||||
|
||||
make_log(
|
||||
"ConvertProcess",
|
||||
f"Processing content {unprocessed_encrypted_content.id} as {content_kind} with preview interval {preview_interval}",
|
||||
level="info"
|
||||
)
|
||||
|
||||
# Выбираем опции конвертации для видео и аудио
|
||||
if content_kind == "video":
|
||||
REQUIRED_CONVERT_OPTIONS = ['high', 'low', 'low_preview']
|
||||
else:
|
||||
REQUIRED_CONVERT_OPTIONS = ['high', 'low'] # no preview for audio
|
||||
|
||||
converted_content = {}
|
||||
logs_dir = "/Storage/logs/converter"
|
||||
|
||||
for option in REQUIRED_CONVERT_OPTIONS:
|
||||
# Set quality parameter and trim option (only for preview)
|
||||
if option == "low_preview":
|
||||
quality = "low"
|
||||
trim_value = f"{preview_interval[0]}-{preview_interval[1]}"
|
||||
else:
|
||||
quality = option
|
||||
trim_value = None
|
||||
|
||||
# Generate a unique output directory for docker container
|
||||
output_uuid = str(uuid.uuid4())
|
||||
output_dir = f"/Storage/storedContent/converter-output/{output_uuid}"
|
||||
|
||||
# Build the docker command
|
||||
cmd = [
|
||||
"docker", "run", "--rm",
|
||||
"-v", f"{input_file_path}:/app/input",
|
||||
"-v", f"{output_dir}:/app/output",
|
||||
"-v", f"{logs_dir}:/app/logs",
|
||||
"media_converter",
|
||||
"--ext", input_ext,
|
||||
"--quality", quality
|
||||
logger.info("Starting media conversion service")
|
||||
|
||||
# Initialize Redis connection
|
||||
self.redis_client = redis.from_url(
|
||||
self.settings.redis_url,
|
||||
encoding="utf-8",
|
||||
decode_responses=True,
|
||||
socket_keepalive=True,
|
||||
socket_keepalive_options={},
|
||||
health_check_interval=30,
|
||||
)
|
||||
|
||||
# Test Redis connection
|
||||
await self.redis_client.ping()
|
||||
logger.info("Redis connection established for converter")
|
||||
|
||||
# Start conversion tasks
|
||||
self.is_running = True
|
||||
|
||||
# Create conversion tasks
|
||||
tasks = [
|
||||
asyncio.create_task(self._process_pending_files_loop()),
|
||||
asyncio.create_task(self._cleanup_temp_files_loop()),
|
||||
asyncio.create_task(self._retry_failed_conversions_loop()),
|
||||
]
|
||||
if trim_value:
|
||||
cmd.extend(["--trim", trim_value])
|
||||
if content_kind == "audio":
|
||||
cmd.append("--audio-only") # audio-only flag
|
||||
|
||||
process = await asyncio.create_subprocess_exec(
|
||||
*cmd,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE
|
||||
)
|
||||
stdout, stderr = await process.communicate()
|
||||
if process.returncode != 0:
|
||||
make_log("ConvertProcess", f"Docker conversion failed for option {option}: {stderr.decode()}", level="error")
|
||||
return
|
||||
|
||||
# List files in output dir
|
||||
|
||||
self.tasks.update(tasks)
|
||||
|
||||
# Wait for all tasks
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error starting conversion service: {e}")
|
||||
await self.stop()
|
||||
raise
|
||||
|
||||
async def stop(self) -> None:
|
||||
"""Stop the conversion service."""
|
||||
logger.info("Stopping media conversion service")
|
||||
self.is_running = False
|
||||
|
||||
# Cancel all tasks
|
||||
for task in self.tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
|
||||
# Wait for tasks to complete
|
||||
if self.tasks:
|
||||
await asyncio.gather(*self.tasks, return_exceptions=True)
|
||||
|
||||
# Close Redis connection
|
||||
if self.redis_client:
|
||||
await self.redis_client.close()
|
||||
|
||||
# Cleanup temp directory
|
||||
await self._cleanup_temp_directory()
|
||||
|
||||
logger.info("Conversion service stopped")
|
||||
|
||||
async def _process_pending_files_loop(self) -> None:
|
||||
"""Main loop for processing pending file conversions."""
|
||||
logger.info("Starting file conversion loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
files = os.listdir(output_dir.replace("/Storage/storedContent", "/app/data"))
|
||||
await self._process_pending_files()
|
||||
await asyncio.sleep(self.process_interval)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
make_log("ConvertProcess", f"Error reading output directory {output_dir}: {e}", level="error")
|
||||
return
|
||||
|
||||
media_files = [f for f in files if f != "output.json"]
|
||||
if len(media_files) != 1:
|
||||
make_log("ConvertProcess", f"Expected one media file, found {len(media_files)} for option {option}", level="error")
|
||||
return
|
||||
|
||||
output_file = os.path.join(
|
||||
output_dir.replace("/Storage/storedContent", "/app/data"),
|
||||
media_files[0]
|
||||
)
|
||||
|
||||
# Compute SHA256 hash of the output file
|
||||
hash_process = await asyncio.create_subprocess_exec(
|
||||
"sha256sum", output_file,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE
|
||||
)
|
||||
hash_stdout, hash_stderr = await hash_process.communicate()
|
||||
if hash_process.returncode != 0:
|
||||
make_log("ConvertProcess", f"Error computing sha256sum for option {option}: {hash_stderr.decode()}", level="error")
|
||||
return
|
||||
file_hash = hash_stdout.decode().split()[0]
|
||||
file_hash = b58encode(bytes.fromhex(file_hash)).decode()
|
||||
|
||||
# Save new StoredContent if not exists
|
||||
if not session.query(StoredContent).filter(
|
||||
StoredContent.hash == file_hash
|
||||
).first():
|
||||
new_content = StoredContent(
|
||||
type="local/content_bin",
|
||||
hash=file_hash,
|
||||
user_id=unprocessed_encrypted_content.user_id,
|
||||
filename=media_files[0],
|
||||
meta={'encrypted_file_hash': unprocessed_encrypted_content.hash},
|
||||
created=datetime.now(),
|
||||
logger.error(f"Error in file conversion loop: {e}")
|
||||
await asyncio.sleep(self.process_interval)
|
||||
|
||||
async def _process_pending_files(self) -> None:
|
||||
"""Process pending file conversions."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get pending uploads
|
||||
result = await session.execute(
|
||||
select(FileUpload)
|
||||
.where(
|
||||
FileUpload.status == "uploaded",
|
||||
FileUpload.processed == False
|
||||
)
|
||||
.limit(self.batch_size)
|
||||
)
|
||||
session.add(new_content)
|
||||
session.commit()
|
||||
|
||||
save_path = os.path.join(UPLOADS_DIR, file_hash)
|
||||
try:
|
||||
os.remove(save_path)
|
||||
except FileNotFoundError:
|
||||
pass
|
||||
|
||||
try:
|
||||
shutil.move(output_file, save_path)
|
||||
uploads = result.scalars().all()
|
||||
|
||||
if not uploads:
|
||||
return
|
||||
|
||||
logger.info(f"Processing {len(uploads)} pending files")
|
||||
|
||||
# Process each upload
|
||||
for upload in uploads:
|
||||
await self._process_single_file(session, upload)
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
make_log("ConvertProcess", f"Error moving output file {output_file} to {save_path}: {e}", level="error")
|
||||
return
|
||||
|
||||
converted_content[option] = file_hash
|
||||
|
||||
# Process output.json for ffprobe_meta
|
||||
output_json_path = os.path.join(
|
||||
output_dir.replace("/Storage/storedContent", "/app/data"),
|
||||
"output.json"
|
||||
)
|
||||
if os.path.exists(output_json_path) and unprocessed_encrypted_content.meta.get('ffprobe_meta') is None:
|
||||
try:
|
||||
with open(output_json_path, "r") as f:
|
||||
ffprobe_meta = json.load(f)
|
||||
unprocessed_encrypted_content.meta = {
|
||||
**unprocessed_encrypted_content.meta,
|
||||
'ffprobe_meta': ffprobe_meta
|
||||
}
|
||||
except Exception as e:
|
||||
make_log("ConvertProcess", f"Error handling output.json for option {option}: {e}", level="error")
|
||||
|
||||
# Cleanup output directory
|
||||
try:
|
||||
shutil.rmtree(output_dir.replace("/Storage/storedContent", "/app/data"))
|
||||
except Exception as e:
|
||||
make_log("ConvertProcess", f"Error removing output dir {output_dir}: {e}", level="warning")
|
||||
|
||||
# Finalize original record
|
||||
make_log("ConvertProcess", f"Content {unprocessed_encrypted_content.id} processed. Converted content: {converted_content}", level="info")
|
||||
unprocessed_encrypted_content.btfs_cid = ContentId(
|
||||
version=2, content_hash=b58decode(converted_content['high' if content_kind=='video' else 'low'])
|
||||
).serialize_v2()
|
||||
unprocessed_encrypted_content.ipfs_cid = ContentId(
|
||||
version=2, content_hash=b58decode(converted_content['low'])
|
||||
).serialize_v2()
|
||||
unprocessed_encrypted_content.meta = {
|
||||
**unprocessed_encrypted_content.meta,
|
||||
'converted_content': converted_content
|
||||
}
|
||||
session.commit()
|
||||
|
||||
# Notify user if needed
|
||||
if not unprocessed_encrypted_content.meta.get('upload_notify_msg_id'):
|
||||
wallet_owner_connection = session.query(WalletConnection).filter(
|
||||
WalletConnection.wallet_address == unprocessed_encrypted_content.owner_address
|
||||
).order_by(WalletConnection.id.desc()).first()
|
||||
if wallet_owner_connection:
|
||||
wallet_owner_user = wallet_owner_connection.user
|
||||
bot = Wrapped_CBotChat(
|
||||
memory._client_telegram_bot,
|
||||
chat_id=wallet_owner_user.telegram_id,
|
||||
user=wallet_owner_user,
|
||||
db_session=session
|
||||
)
|
||||
unprocessed_encrypted_content.meta['upload_notify_msg_id'] = await bot.send_content(session, unprocessed_encrypted_content)
|
||||
session.commit()
|
||||
|
||||
|
||||
async def main_fn(memory):
|
||||
make_log("ConvertProcess", "Service started", level="info")
|
||||
seqno = 0
|
||||
while True:
|
||||
logger.error(f"Error processing pending files: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _process_single_file(self, session: AsyncSession, upload: FileUpload) -> None:
|
||||
"""Process a single file upload."""
|
||||
try:
|
||||
make_log("ConvertProcess", "Service running", level="debug")
|
||||
await convert_loop(memory)
|
||||
await asyncio.sleep(5)
|
||||
await send_status("convert_service", f"working (seqno={seqno})")
|
||||
seqno += 1
|
||||
except BaseException as e:
|
||||
make_log("ConvertProcess", f"Error: {e}", level="error")
|
||||
await asyncio.sleep(3)
|
||||
logger.info(f"Processing file: {upload.filename}")
|
||||
|
||||
# Mark as processing
|
||||
upload.status = "processing"
|
||||
upload.processing_started_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
# Get file extension
|
||||
file_ext = Path(upload.filename).suffix.lower()
|
||||
|
||||
# Process based on file type
|
||||
if file_ext in self.image_formats:
|
||||
await self._process_image(session, upload)
|
||||
elif file_ext in self.video_formats:
|
||||
await self._process_video(session, upload)
|
||||
elif file_ext in self.audio_formats:
|
||||
await self._process_audio(session, upload)
|
||||
elif file_ext in self.document_formats:
|
||||
await self._process_document(session, upload)
|
||||
else:
|
||||
# Just mark as processed for unsupported formats
|
||||
upload.status = "completed"
|
||||
upload.processed = True
|
||||
upload.processing_completed_at = datetime.utcnow()
|
||||
|
||||
# Cache processing result
|
||||
cache_key = f"processed:{upload.id}"
|
||||
processing_info = {
|
||||
"status": upload.status,
|
||||
"processed_at": datetime.utcnow().isoformat(),
|
||||
"metadata": upload.metadata or {}
|
||||
}
|
||||
await self.redis_client.setex(cache_key, 3600, json.dumps(processing_info))
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing file {upload.filename}: {e}")
|
||||
|
||||
# Mark as failed
|
||||
upload.status = "failed"
|
||||
upload.error_message = str(e)
|
||||
upload.retry_count = (upload.retry_count or 0) + 1
|
||||
|
||||
if upload.retry_count >= self.max_retries:
|
||||
upload.processed = True # Stop retrying
|
||||
|
||||
async def _process_image(self, session: AsyncSession, upload: FileUpload) -> None:
|
||||
"""Process an image file."""
|
||||
try:
|
||||
# Download original file
|
||||
original_path = await self._download_file(upload)
|
||||
|
||||
if not original_path:
|
||||
raise Exception("Failed to download original file")
|
||||
|
||||
# Open image
|
||||
with Image.open(original_path) as img:
|
||||
# Extract metadata
|
||||
metadata = {
|
||||
"format": img.format,
|
||||
"mode": img.mode,
|
||||
"size": img.size,
|
||||
"has_transparency": img.mode in ('RGBA', 'LA') or 'transparency' in img.info
|
||||
}
|
||||
|
||||
# Fix orientation
|
||||
img = ImageOps.exif_transpose(img)
|
||||
|
||||
# Resize if too large
|
||||
if img.size[0] > self.max_image_size[0] or img.size[1] > self.max_image_size[1]:
|
||||
img.thumbnail(self.max_image_size, Image.Resampling.LANCZOS)
|
||||
metadata["resized"] = True
|
||||
|
||||
# Save optimized version
|
||||
optimized_path = self.temp_dir / f"optimized_{upload.id}.jpg"
|
||||
|
||||
# Convert to RGB if necessary
|
||||
if img.mode in ('RGBA', 'LA'):
|
||||
background = Image.new('RGB', img.size, (255, 255, 255))
|
||||
if img.mode == 'LA':
|
||||
img = img.convert('RGBA')
|
||||
background.paste(img, mask=img.split()[-1])
|
||||
img = background
|
||||
elif img.mode != 'RGB':
|
||||
img = img.convert('RGB')
|
||||
|
||||
img.save(
|
||||
optimized_path,
|
||||
'JPEG',
|
||||
quality=self.image_quality,
|
||||
optimize=True
|
||||
)
|
||||
|
||||
# Upload optimized version
|
||||
optimized_url = await storage_manager.upload_file(
|
||||
str(optimized_path),
|
||||
f"optimized/{upload.id}/image.jpg"
|
||||
)
|
||||
|
||||
# Generate thumbnails
|
||||
thumbnails = {}
|
||||
for size in self.thumbnail_sizes:
|
||||
thumbnail_path = await self._create_thumbnail(original_path, size)
|
||||
if thumbnail_path:
|
||||
thumb_url = await storage_manager.upload_file(
|
||||
str(thumbnail_path),
|
||||
f"thumbnails/{upload.id}/{size[0]}x{size[1]}.jpg"
|
||||
)
|
||||
thumbnails[f"{size[0]}x{size[1]}"] = thumb_url
|
||||
thumbnail_path.unlink() # Cleanup
|
||||
|
||||
# Update upload record
|
||||
upload.metadata = {
|
||||
**metadata,
|
||||
"thumbnails": thumbnails,
|
||||
"optimized_url": optimized_url
|
||||
}
|
||||
upload.status = "completed"
|
||||
upload.processed = True
|
||||
upload.processing_completed_at = datetime.utcnow()
|
||||
|
||||
# Cleanup temp files
|
||||
original_path.unlink()
|
||||
optimized_path.unlink()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing image {upload.filename}: {e}")
|
||||
raise
|
||||
|
||||
async def _process_video(self, session: AsyncSession, upload: FileUpload) -> None:
|
||||
"""Process a video file."""
|
||||
try:
|
||||
# For video processing, we would typically use ffmpeg
|
||||
# This is a simplified version that just extracts basic info
|
||||
|
||||
original_path = await self._download_file(upload)
|
||||
if not original_path:
|
||||
raise Exception("Failed to download original file")
|
||||
|
||||
# Basic video metadata (would use ffprobe in real implementation)
|
||||
metadata = {
|
||||
"type": "video",
|
||||
"file_size": original_path.stat().st_size,
|
||||
"processing_note": "Video processing requires ffmpeg implementation"
|
||||
}
|
||||
|
||||
# Generate video thumbnail (simplified)
|
||||
thumbnail_path = await self._create_video_thumbnail(original_path)
|
||||
if thumbnail_path:
|
||||
thumb_url = await storage_manager.upload_file(
|
||||
str(thumbnail_path),
|
||||
f"thumbnails/{upload.id}/video_thumb.jpg"
|
||||
)
|
||||
metadata["thumbnail"] = thumb_url
|
||||
thumbnail_path.unlink()
|
||||
|
||||
# Update upload record
|
||||
upload.metadata = metadata
|
||||
upload.status = "completed"
|
||||
upload.processed = True
|
||||
upload.processing_completed_at = datetime.utcnow()
|
||||
|
||||
# Cleanup
|
||||
original_path.unlink()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing video {upload.filename}: {e}")
|
||||
raise
|
||||
|
||||
async def _process_audio(self, session: AsyncSession, upload: FileUpload) -> None:
|
||||
"""Process an audio file."""
|
||||
try:
|
||||
original_path = await self._download_file(upload)
|
||||
if not original_path:
|
||||
raise Exception("Failed to download original file")
|
||||
|
||||
# Basic audio metadata
|
||||
metadata = {
|
||||
"type": "audio",
|
||||
"file_size": original_path.stat().st_size,
|
||||
"processing_note": "Audio processing requires additional libraries"
|
||||
}
|
||||
|
||||
# Update upload record
|
||||
upload.metadata = metadata
|
||||
upload.status = "completed"
|
||||
upload.processed = True
|
||||
upload.processing_completed_at = datetime.utcnow()
|
||||
|
||||
# Cleanup
|
||||
original_path.unlink()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing audio {upload.filename}: {e}")
|
||||
raise
|
||||
|
||||
async def _process_document(self, session: AsyncSession, upload: FileUpload) -> None:
|
||||
"""Process a document file."""
|
||||
try:
|
||||
original_path = await self._download_file(upload)
|
||||
if not original_path:
|
||||
raise Exception("Failed to download original file")
|
||||
|
||||
# Basic document metadata
|
||||
metadata = {
|
||||
"type": "document",
|
||||
"file_size": original_path.stat().st_size,
|
||||
"pages": 1, # Would extract actual page count for PDFs
|
||||
"processing_note": "Document processing requires additional libraries"
|
||||
}
|
||||
|
||||
# Update upload record
|
||||
upload.metadata = metadata
|
||||
upload.status = "completed"
|
||||
upload.processed = True
|
||||
upload.processing_completed_at = datetime.utcnow()
|
||||
|
||||
# Cleanup
|
||||
original_path.unlink()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing document {upload.filename}: {e}")
|
||||
raise
|
||||
|
||||
async def _download_file(self, upload: FileUpload) -> Optional[Path]:
|
||||
"""Download a file for processing."""
|
||||
try:
|
||||
if not upload.file_path:
|
||||
return None
|
||||
|
||||
# Create temp file path
|
||||
temp_path = self.temp_dir / f"original_{upload.id}_{upload.filename}"
|
||||
|
||||
# Download file from storage
|
||||
file_data = await storage_manager.get_file(upload.file_path)
|
||||
if not file_data:
|
||||
return None
|
||||
|
||||
# Write to temp file
|
||||
async with aiofiles.open(temp_path, 'wb') as f:
|
||||
await f.write(file_data)
|
||||
|
||||
return temp_path
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error downloading file {upload.filename}: {e}")
|
||||
return None
|
||||
|
||||
async def _create_thumbnail(self, image_path: Path, size: Tuple[int, int]) -> Optional[Path]:
|
||||
"""Create a thumbnail from an image."""
|
||||
try:
|
||||
thumbnail_path = self.temp_dir / f"thumb_{size[0]}x{size[1]}_{image_path.name}"
|
||||
|
||||
with Image.open(image_path) as img:
|
||||
# Fix orientation
|
||||
img = ImageOps.exif_transpose(img)
|
||||
|
||||
# Create thumbnail
|
||||
img.thumbnail(size, Image.Resampling.LANCZOS)
|
||||
|
||||
# Convert to RGB if necessary
|
||||
if img.mode in ('RGBA', 'LA'):
|
||||
background = Image.new('RGB', img.size, (255, 255, 255))
|
||||
if img.mode == 'LA':
|
||||
img = img.convert('RGBA')
|
||||
background.paste(img, mask=img.split()[-1])
|
||||
img = background
|
||||
elif img.mode != 'RGB':
|
||||
img = img.convert('RGB')
|
||||
|
||||
# Save thumbnail
|
||||
img.save(
|
||||
thumbnail_path,
|
||||
'JPEG',
|
||||
quality=self.image_quality,
|
||||
optimize=True
|
||||
)
|
||||
|
||||
return thumbnail_path
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error creating thumbnail: {e}")
|
||||
return None
|
||||
|
||||
async def _create_video_thumbnail(self, video_path: Path) -> Optional[Path]:
|
||||
"""Create a thumbnail from a video file."""
|
||||
try:
|
||||
# This would require ffmpeg to extract a frame from the video
|
||||
# For now, return a placeholder
|
||||
return None
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error creating video thumbnail: {e}")
|
||||
return None
|
||||
|
||||
async def _cleanup_temp_files_loop(self) -> None:
|
||||
"""Loop for cleaning up temporary files."""
|
||||
logger.info("Starting temp file cleanup loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._cleanup_old_temp_files()
|
||||
await asyncio.sleep(3600) # Run every hour
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in temp cleanup loop: {e}")
|
||||
await asyncio.sleep(3600)
|
||||
|
||||
async def _cleanup_old_temp_files(self) -> None:
|
||||
"""Clean up old temporary files."""
|
||||
try:
|
||||
current_time = datetime.now().timestamp()
|
||||
|
||||
for file_path in self.temp_dir.glob("*"):
|
||||
if file_path.is_file():
|
||||
# Remove files older than 1 hour
|
||||
if current_time - file_path.stat().st_mtime > 3600:
|
||||
file_path.unlink()
|
||||
logger.debug(f"Removed old temp file: {file_path}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error cleaning up temp files: {e}")
|
||||
|
||||
async def _cleanup_temp_directory(self) -> None:
|
||||
"""Clean up the entire temp directory."""
|
||||
try:
|
||||
for file_path in self.temp_dir.glob("*"):
|
||||
if file_path.is_file():
|
||||
file_path.unlink()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error cleaning up temp directory: {e}")
|
||||
|
||||
async def _retry_failed_conversions_loop(self) -> None:
|
||||
"""Loop for retrying failed conversions."""
|
||||
logger.info("Starting retry loop for failed conversions")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._retry_failed_conversions()
|
||||
await asyncio.sleep(1800) # Run every 30 minutes
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in retry loop: {e}")
|
||||
await asyncio.sleep(1800)
|
||||
|
||||
async def _retry_failed_conversions(self) -> None:
|
||||
"""Retry failed conversions that haven't exceeded max retries."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get failed uploads that can be retried
|
||||
result = await session.execute(
|
||||
select(FileUpload)
|
||||
.where(
|
||||
FileUpload.status == "failed",
|
||||
FileUpload.processed == False,
|
||||
(FileUpload.retry_count < self.max_retries) | (FileUpload.retry_count.is_(None))
|
||||
)
|
||||
.limit(5) # Smaller batch for retries
|
||||
)
|
||||
uploads = result.scalars().all()
|
||||
|
||||
for upload in uploads:
|
||||
logger.info(f"Retrying failed conversion for: {upload.filename}")
|
||||
|
||||
# Reset status
|
||||
upload.status = "uploaded"
|
||||
upload.error_message = None
|
||||
|
||||
# Process the file
|
||||
await self._process_single_file(session, upload)
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error retrying failed conversions: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def queue_file_for_processing(self, upload_id: str) -> bool:
|
||||
"""Queue a file for processing."""
|
||||
try:
|
||||
# Add to processing queue
|
||||
queue_key = "conversion_queue"
|
||||
await self.redis_client.lpush(queue_key, upload_id)
|
||||
|
||||
logger.info(f"Queued file {upload_id} for processing")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error queuing file for processing: {e}")
|
||||
return False
|
||||
|
||||
async def get_processing_stats(self) -> Dict[str, Any]:
|
||||
"""Get processing statistics."""
|
||||
try:
|
||||
async with get_async_session() as session:
|
||||
# Get upload stats by status
|
||||
status_result = await session.execute(
|
||||
select(FileUpload.status, asyncio.func.count())
|
||||
.group_by(FileUpload.status)
|
||||
)
|
||||
status_stats = dict(status_result.fetchall())
|
||||
|
||||
# Get processing stats
|
||||
processed_result = await session.execute(
|
||||
select(asyncio.func.count())
|
||||
.select_from(FileUpload)
|
||||
.where(FileUpload.processed == True)
|
||||
)
|
||||
processed_count = processed_result.scalar()
|
||||
|
||||
# Get failed stats
|
||||
failed_result = await session.execute(
|
||||
select(asyncio.func.count())
|
||||
.select_from(FileUpload)
|
||||
.where(FileUpload.status == "failed")
|
||||
)
|
||||
failed_count = failed_result.scalar()
|
||||
|
||||
return {
|
||||
"status_stats": status_stats,
|
||||
"processed_count": processed_count,
|
||||
"failed_count": failed_count,
|
||||
"is_running": self.is_running,
|
||||
"active_tasks": len([t for t in self.tasks if not t.done()]),
|
||||
"temp_files": len(list(self.temp_dir.glob("*"))),
|
||||
"last_update": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting processing stats: {e}")
|
||||
return {"error": str(e)}
|
||||
|
||||
|
||||
# Global converter instance
|
||||
convert_service = ConvertService()
|
||||
@@ -1,313 +1,500 @@
|
||||
"""Blockchain indexer service for monitoring transactions and events."""
|
||||
|
||||
import asyncio
|
||||
from base64 import b64decode
|
||||
from datetime import datetime
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Dict, List, Optional, Set, Any
|
||||
|
||||
from base58 import b58encode
|
||||
from sqlalchemy import String, and_, desc, cast
|
||||
from tonsdk.boc import Cell
|
||||
from tonsdk.utils import Address
|
||||
from app.core._config import CLIENT_TELEGRAM_BOT_USERNAME
|
||||
from app.core._blockchain.ton.platform import platform
|
||||
from app.core._blockchain.ton.toncenter import toncenter
|
||||
from app.core._utils.send_status import send_status
|
||||
from app.core.logger import make_log
|
||||
from app.core.models import UserContent, KnownTelegramMessage, ServiceConfig
|
||||
from app.core.models.node_storage import StoredContent
|
||||
from app.core._utils.resolve_content import resolve_content
|
||||
from app.core.models.wallet_connection import WalletConnection
|
||||
from app.core._keyboards import get_inline_keyboard
|
||||
from app.core.models._telegram import Wrapped_CBotChat
|
||||
from app.core.storage import db_session
|
||||
import os
|
||||
import traceback
|
||||
import redis.asyncio as redis
|
||||
from sqlalchemy import select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.core.database import get_async_session
|
||||
from app.core.models.blockchain import Transaction, Wallet, BlockchainNFT, BlockchainTokenBalance
|
||||
from app.core.background.ton_service import TONService
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
|
||||
if not platform_found:
|
||||
platform_state = await toncenter.get_account(platform.address.to_string(1, 1, 1))
|
||||
if not platform_state.get('code'):
|
||||
make_log("TON", "Platform contract is not deployed, skipping loop", level="info")
|
||||
await send_status("indexer", "not working: platform is not deployed")
|
||||
return False, seqno
|
||||
else:
|
||||
platform_found = True
|
||||
class IndexerService:
|
||||
"""Service for indexing blockchain transactions and events."""
|
||||
|
||||
make_log("Indexer", "Service running", level="debug")
|
||||
with db_session() as session:
|
||||
def __init__(self):
|
||||
self.settings = get_settings()
|
||||
self.ton_service = TONService()
|
||||
self.redis_client: Optional[redis.Redis] = None
|
||||
self.is_running = False
|
||||
self.tasks: Set[asyncio.Task] = set()
|
||||
|
||||
# Indexing configuration
|
||||
self.batch_size = 100
|
||||
self.index_interval = 30 # seconds
|
||||
self.confirmation_blocks = 12
|
||||
self.max_retries = 3
|
||||
|
||||
async def start(self) -> None:
|
||||
"""Start the indexer service."""
|
||||
try:
|
||||
result = await toncenter.run_get_method('EQD8TJ8xEWB1SpnRE4d89YO3jl0W0EiBnNS4IBaHaUmdfizE', 'get_pool_data')
|
||||
assert result['exit_code'] == 0, f"Error in get-method: {result}"
|
||||
assert result['stack'][0][0] == 'num', f"get first element is not num"
|
||||
assert result['stack'][1][0] == 'num', f"get second element is not num"
|
||||
usdt_per_ton = (int(result['stack'][0][1], 16) * 1e3) / int(result['stack'][1][1], 16)
|
||||
ton_per_star = 0.014 / usdt_per_ton
|
||||
ServiceConfig(session).set('live_tonPerStar', [ton_per_star, datetime.utcnow().timestamp()])
|
||||
make_log("TON_Daemon", f"TON per STAR price: {ton_per_star}", level="DEBUG")
|
||||
except BaseException as e:
|
||||
make_log("TON_Daemon", f"Error while saving TON per STAR price: {e}" + '\n' + traceback.format_exc(), level="ERROR")
|
||||
|
||||
new_licenses = session.query(UserContent).filter(
|
||||
and_(
|
||||
~UserContent.meta.contains({'notification_sent': True}),
|
||||
UserContent.type == 'nft/listen'
|
||||
logger.info("Starting blockchain indexer service")
|
||||
|
||||
# Initialize Redis connection
|
||||
self.redis_client = redis.from_url(
|
||||
self.settings.redis_url,
|
||||
encoding="utf-8",
|
||||
decode_responses=True,
|
||||
socket_keepalive=True,
|
||||
socket_keepalive_options={},
|
||||
health_check_interval=30,
|
||||
)
|
||||
).all()
|
||||
for new_license in new_licenses:
|
||||
licensed_content = session.query(StoredContent).filter(
|
||||
StoredContent.id == new_license.content_id
|
||||
).first()
|
||||
if not licensed_content:
|
||||
make_log("Indexer", f"Licensed content not found: {new_license.content_id}", level="error")
|
||||
|
||||
content_metadata = licensed_content.metadata_json(session)
|
||||
assert content_metadata, "No content metadata found"
|
||||
|
||||
if not (licensed_content.owner_address == new_license.owner_address):
|
||||
try:
|
||||
user = new_license.user
|
||||
if user.telegram_id and licensed_content:
|
||||
await (Wrapped_CBotChat(memory._client_telegram_bot, chat_id=user.telegram_id, user=user, db_session=session)).send_content(
|
||||
session, licensed_content
|
||||
)
|
||||
|
||||
wallet_owner_connection = session.query(WalletConnection).filter_by(
|
||||
wallet_address=licensed_content.owner_address,
|
||||
invalidated=False
|
||||
).order_by(desc(WalletConnection.id)).first()
|
||||
wallet_owner_user = wallet_owner_connection.user
|
||||
if wallet_owner_user.telegram_id:
|
||||
wallet_owner_bot = Wrapped_CBotChat(memory._telegram_bot, chat_id=wallet_owner_user.telegram_id, user=wallet_owner_user, db_session=session)
|
||||
await wallet_owner_bot.send_message(
|
||||
user.translated('p_licenseWasBought').format(
|
||||
username=user.front_format(),
|
||||
nft_address=f'"https://tonviewer.com/{new_license.onchain_address}"',
|
||||
content_title=content_metadata.get('name', 'Unknown'),
|
||||
),
|
||||
message_type='notification',
|
||||
)
|
||||
except BaseException as e:
|
||||
make_log("IndexerSendNewLicense", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
|
||||
|
||||
new_license.meta = {**new_license.meta, 'notification_sent': True}
|
||||
session.commit()
|
||||
|
||||
content_without_cid = session.query(StoredContent).filter(
|
||||
StoredContent.content_id == None
|
||||
)
|
||||
for target_content in content_without_cid:
|
||||
target_cid = target_content.cid.serialize_v2()
|
||||
make_log("Indexer", f"Content without CID: {target_content.hash}, setting CID: {target_cid}", level="debug")
|
||||
target_content.content_id = target_cid
|
||||
|
||||
session.commit()
|
||||
|
||||
last_known_index_ = session.query(StoredContent).filter(
|
||||
StoredContent.onchain_index != None
|
||||
).order_by(StoredContent.onchain_index.desc()).first()
|
||||
last_known_index = last_known_index_.onchain_index if last_known_index_ else 0
|
||||
last_known_index = max(last_known_index, 0)
|
||||
make_log("Indexer", f"Last known index: {last_known_index}", level="debug")
|
||||
if last_known_index_:
|
||||
next_item_index = last_known_index + 1
|
||||
else:
|
||||
next_item_index = 0
|
||||
|
||||
resolve_item_result = await toncenter.run_get_method(platform.address.to_string(1, 1, 1), 'get_nft_address_by_index', [['num', next_item_index]])
|
||||
make_log("Indexer", f"Resolve item result: {resolve_item_result}", level="debug")
|
||||
if resolve_item_result.get('exit_code', -1) != 0:
|
||||
make_log("Indexer", f"Resolve item error: {resolve_item_result}", level="error")
|
||||
return platform_found, seqno
|
||||
|
||||
item_address_cell_b64 = resolve_item_result['stack'][0][1]["bytes"]
|
||||
item_address_slice = Cell.one_from_boc(b64decode(item_address_cell_b64)).begin_parse()
|
||||
item_address = item_address_slice.read_msg_addr()
|
||||
make_log("Indexer", f"Item address: {item_address.to_string(1, 1, 1)}", level="debug")
|
||||
|
||||
item_get_data_result = await toncenter.run_get_method(item_address.to_string(1, 1, 1), 'indexator_data')
|
||||
if item_get_data_result.get('exit_code', -1) != 0:
|
||||
make_log("Indexer", f"Get item data error (maybe not deployed): {item_get_data_result}", level="debug")
|
||||
return platform_found, seqno
|
||||
|
||||
assert item_get_data_result['stack'][0][0] == 'num', "Item type is not a number"
|
||||
assert int(item_get_data_result['stack'][0][1], 16) == 1, "Item is not COP NFT"
|
||||
item_returned_address = Cell.one_from_boc(b64decode(item_get_data_result['stack'][1][1]['bytes'])).begin_parse().read_msg_addr()
|
||||
assert (
|
||||
item_returned_address.to_string(1, 1, 1) == item_address.to_string(1, 1, 1)
|
||||
), "Item address mismatch"
|
||||
|
||||
assert item_get_data_result['stack'][2][0] == 'num', "Item index is not a number"
|
||||
item_index = int(item_get_data_result['stack'][2][1], 16)
|
||||
assert item_index == next_item_index, "Item index mismatch"
|
||||
|
||||
item_platform_address = Cell.one_from_boc(b64decode(item_get_data_result['stack'][3][1]['bytes'])).begin_parse().read_msg_addr()
|
||||
assert item_platform_address.to_string(1, 1, 1) == Address(platform.address.to_string(1, 1, 1)).to_string(1, 1, 1), "Item platform address mismatch"
|
||||
|
||||
assert item_get_data_result['stack'][4][0] == 'num', "Item license type is not a number"
|
||||
item_license_type = int(item_get_data_result['stack'][4][1], 16)
|
||||
assert item_license_type == 0, "Item license type is not 0"
|
||||
|
||||
item_owner_address = Cell.one_from_boc(b64decode(item_get_data_result['stack'][5][1]["bytes"])).begin_parse().read_msg_addr()
|
||||
item_values = Cell.one_from_boc(b64decode(item_get_data_result['stack'][6][1]['bytes']))
|
||||
item_derivates = Cell.one_from_boc(b64decode(item_get_data_result['stack'][7][1]['bytes']))
|
||||
item_platform_variables = Cell.one_from_boc(b64decode(item_get_data_result['stack'][8][1]['bytes']))
|
||||
item_distribution = Cell.one_from_boc(b64decode(item_get_data_result['stack'][9][1]['bytes']))
|
||||
|
||||
item_distribution_slice = item_distribution.begin_parse()
|
||||
item_prices_slice = item_distribution_slice.refs[0].begin_parse()
|
||||
item_listen_license_price = item_prices_slice.read_coins()
|
||||
item_use_license_price = item_prices_slice.read_coins()
|
||||
item_resale_license_price = item_prices_slice.read_coins()
|
||||
|
||||
item_values_slice = item_values.begin_parse()
|
||||
item_content_hash_int = item_values_slice.read_uint(256)
|
||||
item_content_hash = item_content_hash_int.to_bytes(32, 'big')
|
||||
# item_content_hash_str = b58encode(item_content_hash).decode()
|
||||
item_metadata = item_values_slice.refs[0]
|
||||
item_content = item_values_slice.refs[1]
|
||||
item_metadata_str = item_metadata.bits.array.decode()
|
||||
item_content_cid_str = item_content.refs[0].bits.array.decode()
|
||||
item_content_cover_cid_str = item_content.refs[1].bits.array.decode()
|
||||
item_content_metadata_cid_str = item_content.refs[2].bits.array.decode()
|
||||
|
||||
item_content_cid, err = resolve_content(item_content_cid_str)
|
||||
item_content_hash = item_content_cid.content_hash
|
||||
item_content_hash_str = item_content_cid.content_hash_b58
|
||||
|
||||
item_metadata_packed = {
|
||||
'license_type': item_license_type,
|
||||
'item_address': item_address.to_string(1, 1, 1),
|
||||
'content_cid': item_content_cid_str,
|
||||
'cover_cid': item_content_cover_cid_str,
|
||||
'metadata_cid': item_content_metadata_cid_str,
|
||||
'derivates': b58encode(item_derivates.to_boc(False)).decode(),
|
||||
'platform_variables': b58encode(item_platform_variables.to_boc(False)).decode(),
|
||||
'license': {
|
||||
'listen': {
|
||||
'price': str(item_listen_license_price)
|
||||
},
|
||||
'use': {
|
||||
'price': str(item_use_license_price)
|
||||
},
|
||||
'resale': {
|
||||
'price': str(item_resale_license_price)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
user_wallet_connection = None
|
||||
if item_owner_address:
|
||||
user_wallet_connection = session.query(WalletConnection).filter(
|
||||
WalletConnection.wallet_address == item_owner_address.to_string(1, 1, 1)
|
||||
).first()
|
||||
|
||||
encrypted_stored_content = session.query(StoredContent).filter(
|
||||
StoredContent.hash == item_content_hash_str,
|
||||
# StoredContent.type.like("local%")
|
||||
).first()
|
||||
if encrypted_stored_content:
|
||||
is_duplicate = encrypted_stored_content.type.startswith("onchain") \
|
||||
and encrypted_stored_content.onchain_index != item_index
|
||||
if not is_duplicate:
|
||||
if encrypted_stored_content.type.startswith('local'):
|
||||
encrypted_stored_content.type = "onchain/content" + ("_unknown" if (encrypted_stored_content.key_id is None) else "")
|
||||
encrypted_stored_content.onchain_index = item_index
|
||||
encrypted_stored_content.owner_address = item_owner_address.to_string(1, 1, 1)
|
||||
user = None
|
||||
if user_wallet_connection:
|
||||
encrypted_stored_content.user_id = user_wallet_connection.user_id
|
||||
user = user_wallet_connection.user
|
||||
|
||||
if user:
|
||||
user_uploader_wrapper = Wrapped_CBotChat(memory._telegram_bot, chat_id=user.telegram_id, user=user, db_session=session)
|
||||
await user_uploader_wrapper.send_message(
|
||||
user.translated('p_contentWasIndexed').format(
|
||||
item_address=item_address.to_string(1, 1, 1),
|
||||
item_index=item_index,
|
||||
),
|
||||
message_type='notification',
|
||||
reply_markup=get_inline_keyboard([
|
||||
[{
|
||||
'text': user.translated('viewTrackAsClient_button'),
|
||||
'url': f"https://t.me/{CLIENT_TELEGRAM_BOT_USERNAME}?start=C{encrypted_stored_content.cid.serialize_v2()}"
|
||||
}],
|
||||
])
|
||||
)
|
||||
|
||||
try:
|
||||
for hint_message in session.query(KnownTelegramMessage).filter(
|
||||
and_(
|
||||
KnownTelegramMessage.chat_id == user.telegram_id,
|
||||
KnownTelegramMessage.type == 'hint',
|
||||
cast(KnownTelegramMessage.meta['encrypted_content_hash'], String) == encrypted_stored_content.hash,
|
||||
KnownTelegramMessage.deleted == False
|
||||
)
|
||||
).all():
|
||||
await user_uploader_wrapper.delete_message(hint_message.message_id)
|
||||
except BaseException as e:
|
||||
make_log("Indexer", f"Error while deleting hint messages: {e}" + '\n' + traceback.format_exc(), level="error")
|
||||
elif encrypted_stored_content.type.startswith('onchain') and encrypted_stored_content.onchain_index == item_index:
|
||||
encrypted_stored_content.type = "onchain/content" + ("_unknown" if (encrypted_stored_content.key_id is None) else "")
|
||||
encrypted_stored_content.owner_address = item_owner_address.to_string(1, 1, 1)
|
||||
if user_wallet_connection:
|
||||
encrypted_stored_content.user_id = user_wallet_connection.user_id
|
||||
else:
|
||||
make_log("Indexer", f"[CRITICAL] Item already indexed and ERRORED!: {item_content_hash_str}", level="error")
|
||||
return platform_found, seqno
|
||||
|
||||
encrypted_stored_content.updated = datetime.now()
|
||||
encrypted_stored_content.meta = {
|
||||
**encrypted_stored_content.meta,
|
||||
**item_metadata_packed
|
||||
}
|
||||
|
||||
session.commit()
|
||||
return platform_found, seqno
|
||||
else:
|
||||
item_metadata_packed['copied_from'] = encrypted_stored_content.id
|
||||
item_metadata_packed['copied_from_cid'] = encrypted_stored_content.cid.serialize_v2()
|
||||
item_content_hash_str = f"{b58encode(bytes(16) + os.urandom(30)).decode()}" # check this for vulnerability
|
||||
|
||||
onchain_stored_content = StoredContent(
|
||||
type="onchain/content_unknown",
|
||||
hash=item_content_hash_str,
|
||||
onchain_index=item_index,
|
||||
owner_address=item_owner_address.to_string(1, 1, 1) if item_owner_address else None,
|
||||
meta=item_metadata_packed,
|
||||
filename="UNKNOWN_ENCRYPTED_CONTENT",
|
||||
user_id=user_wallet_connection.user_id if user_wallet_connection else None,
|
||||
created=datetime.now(),
|
||||
encrypted=True,
|
||||
decrypted_content_id=None,
|
||||
key_id=None,
|
||||
updated=datetime.now()
|
||||
)
|
||||
session.add(onchain_stored_content)
|
||||
session.commit()
|
||||
make_log("Indexer", f"Item indexed: {item_content_hash_str}", level="info")
|
||||
last_known_index += 1
|
||||
|
||||
return platform_found, seqno
|
||||
|
||||
|
||||
async def main_fn(memory, ):
|
||||
make_log("Indexer", "Service started", level="info")
|
||||
platform_found = False
|
||||
seqno = 0
|
||||
while True:
|
||||
|
||||
# Test Redis connection
|
||||
await self.redis_client.ping()
|
||||
logger.info("Redis connection established for indexer")
|
||||
|
||||
# Start indexing tasks
|
||||
self.is_running = True
|
||||
|
||||
# Create indexing tasks
|
||||
tasks = [
|
||||
asyncio.create_task(self._index_transactions_loop()),
|
||||
asyncio.create_task(self._index_wallets_loop()),
|
||||
asyncio.create_task(self._index_nfts_loop()),
|
||||
asyncio.create_task(self._index_token_balances_loop()),
|
||||
asyncio.create_task(self._cleanup_cache_loop()),
|
||||
]
|
||||
|
||||
self.tasks.update(tasks)
|
||||
|
||||
# Wait for all tasks
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error starting indexer service: {e}")
|
||||
await self.stop()
|
||||
raise
|
||||
|
||||
async def stop(self) -> None:
|
||||
"""Stop the indexer service."""
|
||||
logger.info("Stopping blockchain indexer service")
|
||||
self.is_running = False
|
||||
|
||||
# Cancel all tasks
|
||||
for task in self.tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
|
||||
# Wait for tasks to complete
|
||||
if self.tasks:
|
||||
await asyncio.gather(*self.tasks, return_exceptions=True)
|
||||
|
||||
# Close Redis connection
|
||||
if self.redis_client:
|
||||
await self.redis_client.close()
|
||||
|
||||
logger.info("Indexer service stopped")
|
||||
|
||||
async def _index_transactions_loop(self) -> None:
|
||||
"""Main loop for indexing transactions."""
|
||||
logger.info("Starting transaction indexing loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._index_pending_transactions()
|
||||
await self._update_transaction_confirmations()
|
||||
await asyncio.sleep(self.index_interval)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in transaction indexing loop: {e}")
|
||||
await asyncio.sleep(self.index_interval)
|
||||
|
||||
async def _index_wallets_loop(self) -> None:
|
||||
"""Loop for updating wallet information."""
|
||||
logger.info("Starting wallet indexing loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._update_wallet_balances()
|
||||
await asyncio.sleep(self.index_interval * 2) # Less frequent
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in wallet indexing loop: {e}")
|
||||
await asyncio.sleep(self.index_interval * 2)
|
||||
|
||||
async def _index_nfts_loop(self) -> None:
|
||||
"""Loop for indexing NFT collections and transfers."""
|
||||
logger.info("Starting NFT indexing loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._index_nft_collections()
|
||||
await self._index_nft_transfers()
|
||||
await asyncio.sleep(self.index_interval * 4) # Even less frequent
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in NFT indexing loop: {e}")
|
||||
await asyncio.sleep(self.index_interval * 4)
|
||||
|
||||
async def _index_token_balances_loop(self) -> None:
|
||||
"""Loop for updating token balances."""
|
||||
logger.info("Starting token balance indexing loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._update_token_balances()
|
||||
await asyncio.sleep(self.index_interval * 3)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in token balance indexing loop: {e}")
|
||||
await asyncio.sleep(self.index_interval * 3)
|
||||
|
||||
async def _cleanup_cache_loop(self) -> None:
|
||||
"""Loop for cleaning up old cache entries."""
|
||||
logger.info("Starting cache cleanup loop")
|
||||
|
||||
while self.is_running:
|
||||
try:
|
||||
await self._cleanup_old_cache_entries()
|
||||
await asyncio.sleep(3600) # Run every hour
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Error in cache cleanup loop: {e}")
|
||||
await asyncio.sleep(3600)
|
||||
|
||||
async def _index_pending_transactions(self) -> None:
|
||||
"""Index pending transactions from the database."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get pending transactions
|
||||
result = await session.execute(
|
||||
select(Transaction)
|
||||
.where(Transaction.status == "pending")
|
||||
.limit(self.batch_size)
|
||||
)
|
||||
transactions = result.scalars().all()
|
||||
|
||||
if not transactions:
|
||||
return
|
||||
|
||||
logger.info(f"Indexing {len(transactions)} pending transactions")
|
||||
|
||||
# Process each transaction
|
||||
for transaction in transactions:
|
||||
await self._process_transaction(session, transaction)
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error indexing pending transactions: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _process_transaction(self, session: AsyncSession, transaction: Transaction) -> None:
|
||||
"""Process a single transaction."""
|
||||
try:
|
||||
platform_found, seqno = await indexer_loop(memory, platform_found, seqno)
|
||||
except BaseException as e:
|
||||
make_log("Indexer", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
|
||||
|
||||
if platform_found:
|
||||
await send_status("indexer", f"working (seqno={seqno})")
|
||||
|
||||
await asyncio.sleep(5)
|
||||
seqno += 1
|
||||
# Check transaction status on blockchain
|
||||
if transaction.tx_hash:
|
||||
tx_info = await self.ton_service.get_transaction_info(transaction.tx_hash)
|
||||
|
||||
if tx_info:
|
||||
# Update transaction with blockchain data
|
||||
transaction.status = tx_info.get("status", "pending")
|
||||
transaction.block_number = tx_info.get("block_number")
|
||||
transaction.gas_used = tx_info.get("gas_used")
|
||||
transaction.gas_price = tx_info.get("gas_price")
|
||||
transaction.confirmations = tx_info.get("confirmations", 0)
|
||||
transaction.updated_at = datetime.utcnow()
|
||||
|
||||
# Cache transaction info
|
||||
cache_key = f"tx:{transaction.tx_hash}"
|
||||
await self.redis_client.setex(
|
||||
cache_key,
|
||||
3600, # 1 hour
|
||||
json.dumps(tx_info)
|
||||
)
|
||||
|
||||
logger.debug(f"Updated transaction {transaction.tx_hash}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing transaction {transaction.id}: {e}")
|
||||
|
||||
async def _update_transaction_confirmations(self) -> None:
|
||||
"""Update confirmation counts for recent transactions."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get recent confirmed transactions
|
||||
cutoff_time = datetime.utcnow() - timedelta(hours=24)
|
||||
result = await session.execute(
|
||||
select(Transaction)
|
||||
.where(
|
||||
Transaction.status == "confirmed",
|
||||
Transaction.confirmations < self.confirmation_blocks,
|
||||
Transaction.updated_at > cutoff_time
|
||||
)
|
||||
.limit(self.batch_size)
|
||||
)
|
||||
transactions = result.scalars().all()
|
||||
|
||||
for transaction in transactions:
|
||||
if transaction.tx_hash:
|
||||
try:
|
||||
confirmations = await self.ton_service.get_transaction_confirmations(
|
||||
transaction.tx_hash
|
||||
)
|
||||
|
||||
if confirmations != transaction.confirmations:
|
||||
transaction.confirmations = confirmations
|
||||
transaction.updated_at = datetime.utcnow()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating confirmations for {transaction.tx_hash}: {e}")
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating transaction confirmations: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _update_wallet_balances(self) -> None:
|
||||
"""Update wallet balances from the blockchain."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get active wallets
|
||||
result = await session.execute(
|
||||
select(Wallet)
|
||||
.where(Wallet.is_active == True)
|
||||
.limit(self.batch_size)
|
||||
)
|
||||
wallets = result.scalars().all()
|
||||
|
||||
for wallet in wallets:
|
||||
try:
|
||||
# Get current balance
|
||||
balance = await self.ton_service.get_wallet_balance(wallet.address)
|
||||
|
||||
if balance != wallet.balance:
|
||||
wallet.balance = balance
|
||||
wallet.updated_at = datetime.utcnow()
|
||||
|
||||
# Cache balance
|
||||
cache_key = f"balance:{wallet.address}"
|
||||
await self.redis_client.setex(cache_key, 300, str(balance)) # 5 minutes
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating balance for wallet {wallet.address}: {e}")
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating wallet balances: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _index_nft_collections(self) -> None:
|
||||
"""Index NFT collections and metadata."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get wallets to check for NFTs
|
||||
result = await session.execute(
|
||||
select(Wallet)
|
||||
.where(Wallet.is_active == True)
|
||||
.limit(self.batch_size // 4) # Smaller batch for NFTs
|
||||
)
|
||||
wallets = result.scalars().all()
|
||||
|
||||
for wallet in wallets:
|
||||
try:
|
||||
# Get NFTs for this wallet
|
||||
nfts = await self.ton_service.get_wallet_nfts(wallet.address)
|
||||
|
||||
for nft_data in nfts:
|
||||
await self._process_nft(session, wallet, nft_data)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error indexing NFTs for wallet {wallet.address}: {e}")
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error indexing NFT collections: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _process_nft(self, session: AsyncSession, wallet: Wallet, nft_data: Dict[str, Any]) -> None:
|
||||
"""Process a single NFT."""
|
||||
try:
|
||||
# Check if NFT exists
|
||||
result = await session.execute(
|
||||
select(BlockchainNFT)
|
||||
.where(
|
||||
BlockchainNFT.token_id == nft_data["token_id"],
|
||||
BlockchainNFT.collection_address == nft_data["collection_address"]
|
||||
)
|
||||
)
|
||||
existing_nft = result.scalar_one_or_none()
|
||||
|
||||
if existing_nft:
|
||||
# Update existing NFT
|
||||
existing_nft.owner_address = wallet.address
|
||||
existing_nft.metadata = nft_data.get("metadata", {})
|
||||
existing_nft.updated_at = datetime.utcnow()
|
||||
else:
|
||||
# Create new NFT
|
||||
new_nft = BlockchainNFT(
|
||||
wallet_id=wallet.id,
|
||||
token_id=nft_data["token_id"],
|
||||
collection_address=nft_data["collection_address"],
|
||||
owner_address=wallet.address,
|
||||
token_uri=nft_data.get("token_uri"),
|
||||
metadata=nft_data.get("metadata", {}),
|
||||
created_at=datetime.utcnow()
|
||||
)
|
||||
session.add(new_nft)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing NFT {nft_data.get('token_id')}: {e}")
|
||||
|
||||
async def _index_nft_transfers(self) -> None:
|
||||
"""Index NFT transfers."""
|
||||
# This would involve checking recent blocks for NFT transfer events
|
||||
# Implementation depends on the specific blockchain's event system
|
||||
pass
|
||||
|
||||
async def _update_token_balances(self) -> None:
|
||||
"""Update token balances for wallets."""
|
||||
async with get_async_session() as session:
|
||||
try:
|
||||
# Get wallets with token balances to update
|
||||
result = await session.execute(
|
||||
select(Wallet)
|
||||
.where(Wallet.is_active == True)
|
||||
.limit(self.batch_size // 2)
|
||||
)
|
||||
wallets = result.scalars().all()
|
||||
|
||||
for wallet in wallets:
|
||||
try:
|
||||
# Get token balances
|
||||
token_balances = await self.ton_service.get_wallet_token_balances(wallet.address)
|
||||
|
||||
for token_data in token_balances:
|
||||
await self._update_token_balance(session, wallet, token_data)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating token balances for {wallet.address}: {e}")
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating token balances: {e}")
|
||||
await session.rollback()
|
||||
|
||||
async def _update_token_balance(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
wallet: Wallet,
|
||||
token_data: Dict[str, Any]
|
||||
) -> None:
|
||||
"""Update a single token balance."""
|
||||
try:
|
||||
# Check if balance record exists
|
||||
result = await session.execute(
|
||||
select(BlockchainTokenBalance)
|
||||
.where(
|
||||
BlockchainTokenBalance.wallet_id == wallet.id,
|
||||
BlockchainTokenBalance.token_address == token_data["token_address"]
|
||||
)
|
||||
)
|
||||
existing_balance = result.scalar_one_or_none()
|
||||
|
||||
if existing_balance:
|
||||
# Update existing balance
|
||||
existing_balance.balance = token_data["balance"]
|
||||
existing_balance.decimals = token_data.get("decimals", 18)
|
||||
existing_balance.updated_at = datetime.utcnow()
|
||||
else:
|
||||
# Create new balance record
|
||||
new_balance = BlockchainTokenBalance(
|
||||
wallet_id=wallet.id,
|
||||
token_address=token_data["token_address"],
|
||||
token_name=token_data.get("name"),
|
||||
token_symbol=token_data.get("symbol"),
|
||||
balance=token_data["balance"],
|
||||
decimals=token_data.get("decimals", 18),
|
||||
created_at=datetime.utcnow()
|
||||
)
|
||||
session.add(new_balance)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating token balance: {e}")
|
||||
|
||||
async def _cleanup_old_cache_entries(self) -> None:
|
||||
"""Clean up old cache entries."""
|
||||
try:
|
||||
# Get all keys with our prefixes
|
||||
patterns = ["tx:*", "balance:*", "nft:*", "token:*"]
|
||||
|
||||
for pattern in patterns:
|
||||
keys = await self.redis_client.keys(pattern)
|
||||
|
||||
# Check TTL and remove expired keys
|
||||
for key in keys:
|
||||
ttl = await self.redis_client.ttl(key)
|
||||
if ttl == -1: # No expiration set
|
||||
await self.redis_client.expire(key, 3600) # Set 1 hour expiration
|
||||
|
||||
logger.debug("Cache cleanup completed")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error during cache cleanup: {e}")
|
||||
|
||||
async def get_indexing_stats(self) -> Dict[str, Any]:
|
||||
"""Get indexing statistics."""
|
||||
try:
|
||||
async with get_async_session() as session:
|
||||
# Get transaction stats
|
||||
tx_result = await session.execute(
|
||||
select(Transaction.status, asyncio.func.count())
|
||||
.group_by(Transaction.status)
|
||||
)
|
||||
tx_stats = dict(tx_result.fetchall())
|
||||
|
||||
# Get wallet stats
|
||||
wallet_result = await session.execute(
|
||||
select(asyncio.func.count())
|
||||
.select_from(Wallet)
|
||||
.where(Wallet.is_active == True)
|
||||
)
|
||||
active_wallets = wallet_result.scalar()
|
||||
|
||||
# Get NFT stats
|
||||
nft_result = await session.execute(
|
||||
select(asyncio.func.count())
|
||||
.select_from(BlockchainNFT)
|
||||
)
|
||||
total_nfts = nft_result.scalar()
|
||||
|
||||
return {
|
||||
"transaction_stats": tx_stats,
|
||||
"active_wallets": active_wallets,
|
||||
"total_nfts": total_nfts,
|
||||
"is_running": self.is_running,
|
||||
"active_tasks": len([t for t in self.tasks if not t.done()]),
|
||||
"last_update": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting indexing stats: {e}")
|
||||
return {"error": str(e)}
|
||||
|
||||
|
||||
|
||||
# if __name__ == '__main__':
|
||||
# loop = asyncio.get_event_loop()
|
||||
# loop.run_until_complete(main())
|
||||
# loop.close()
|
||||
# Global indexer instance
|
||||
indexer_service = IndexerService()
|
||||
+644
-276
@@ -1,290 +1,658 @@
|
||||
"""
|
||||
TON Blockchain service for wallet operations, transaction management, and smart contract interactions.
|
||||
Provides async operations with connection pooling, caching, and comprehensive error handling.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
from base64 import b64decode
|
||||
import os
|
||||
import traceback
|
||||
import httpx
|
||||
|
||||
from sqlalchemy import and_, func
|
||||
from tonsdk.boc import begin_cell, Cell
|
||||
from tonsdk.contract.wallet import Wallets
|
||||
from tonsdk.utils import HighloadQueryId
|
||||
import json
|
||||
from datetime import datetime, timedelta
|
||||
from decimal import Decimal
|
||||
from typing import Dict, List, Optional, Any, Tuple
|
||||
from uuid import UUID
|
||||
|
||||
from app.core._blockchain.ton.platform import platform
|
||||
from app.core._blockchain.ton.toncenter import toncenter
|
||||
from app.core.models.tasks import BlockchainTask
|
||||
from app.core._config import MY_FUND_ADDRESS
|
||||
from app.core._secrets import service_wallet
|
||||
from app.core._utils.send_status import send_status
|
||||
from app.core.storage import db_session
|
||||
from app.core.logger import make_log
|
||||
import httpx
|
||||
from sqlalchemy import select, update, and_
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.core.database import get_async_session, get_cache_manager
|
||||
from app.core.logging import get_logger
|
||||
from app.core.security import decrypt_data, encrypt_data
|
||||
|
||||
async def get_sw_seqno():
|
||||
sw_seqno_result = await toncenter.run_get_method(service_wallet.address.to_string(1, 1, 1), 'seqno')
|
||||
if sw_seqno_result.get('exit_code', -1) != 0:
|
||||
sw_seqno_value = 0
|
||||
else:
|
||||
sw_seqno_value = int(sw_seqno_result.get('stack', [['num', '0x0']])[0][1], 16)
|
||||
logger = get_logger(__name__)
|
||||
settings = get_settings()
|
||||
|
||||
return sw_seqno_value
|
||||
|
||||
|
||||
async def main_fn(memory):
|
||||
make_log("TON", f"Service started, SW = {service_wallet.address.to_string(1, 1, 1)}", level="info")
|
||||
sw_seqno_value = await get_sw_seqno()
|
||||
make_log("TON", f"Service wallet run seqno method: {sw_seqno_value}", level="info")
|
||||
if sw_seqno_value == 0:
|
||||
make_log("TON", "Service wallet is not deployed, deploying...", level="info")
|
||||
await toncenter.send_boc(
|
||||
service_wallet.create_transfer_message(
|
||||
[{
|
||||
'address': service_wallet.address.to_string(1, 1, 1),
|
||||
'amount': 1,
|
||||
'send_mode': 1,
|
||||
'payload': begin_cell().store_uint(0, 32).store_bytes(b"Init MY Node").end_cell()
|
||||
}], 0
|
||||
)['message'].to_boc(False)
|
||||
)
|
||||
await asyncio.sleep(5)
|
||||
return await main_fn(memory)
|
||||
|
||||
if os.getenv("TON_BEGIN_COMMAND_WITHDRAW"):
|
||||
await toncenter.send_boc(
|
||||
service_wallet.create_transfer_message(
|
||||
[{
|
||||
'address': MY_FUND_ADDRESS,
|
||||
'amount': 1,
|
||||
'send_mode': 128,
|
||||
'payload': begin_cell().end_cell()
|
||||
}], sw_seqno_value
|
||||
)['message'].to_boc(False)
|
||||
)
|
||||
make_log("TON", "Withdraw command sent", level="info")
|
||||
await asyncio.sleep(10)
|
||||
return await main_fn(memory)
|
||||
|
||||
# TODO: не деплоить если указан master_address и мы проверили что аккаунт существует. Сейчас platform у каждой ноды будет разным
|
||||
platform_state = await toncenter.get_account(platform.address.to_string(1, 1, 1))
|
||||
if not platform_state.get('code'):
|
||||
make_log("TON", "Platform contract is not deployed, send deploy transaction..", level="info")
|
||||
await toncenter.send_boc(
|
||||
service_wallet.create_transfer_message(
|
||||
[{
|
||||
'address': platform.address.to_string(1, 1, 1),
|
||||
'amount': int(0.08 * 10 ** 9),
|
||||
'send_mode': 1,
|
||||
'payload': begin_cell().store_uint(0, 32).store_uint(0, 64).end_cell(),
|
||||
'state_init': platform.create_state_init()['state_init']
|
||||
}], sw_seqno_value
|
||||
)['message'].to_boc(False)
|
||||
)
|
||||
|
||||
await send_status("ton_daemon", "working: deploying platform")
|
||||
await asyncio.sleep(15)
|
||||
return await main_fn(memory)
|
||||
class TONService:
|
||||
"""
|
||||
Comprehensive TON blockchain service with async operations.
|
||||
Handles wallet management, transactions, and smart contract interactions.
|
||||
"""
|
||||
|
||||
highload_wallet = Wallets.ALL['hv3'](
|
||||
private_key=service_wallet.options['private_key'],
|
||||
public_key=service_wallet.options['public_key'],
|
||||
wc=0
|
||||
)
|
||||
make_log("TON", f"Highload wallet address: {highload_wallet.address.to_string(1, 1, 1)}", level="info")
|
||||
highload_state = await toncenter.get_account(highload_wallet.address.to_string(1, 1, 1))
|
||||
if int(highload_state.get('balance', '0')) / 1e9 < 0.05:
|
||||
make_log("TON", "Highload wallet balance is less than 0.05, send topup transaction..", level="info")
|
||||
await toncenter.send_boc(
|
||||
service_wallet.create_transfer_message(
|
||||
[{
|
||||
'address': highload_wallet.address.to_string(1, 1, 0),
|
||||
'amount': int(0.08 * 10 ** 9),
|
||||
'send_mode': 1,
|
||||
'payload': begin_cell().store_uint(0, 32).end_cell()
|
||||
}], sw_seqno_value
|
||||
)['message'].to_boc(False)
|
||||
def __init__(self):
|
||||
self.api_endpoint = settings.TON_API_ENDPOINT
|
||||
self.testnet = settings.TON_TESTNET
|
||||
self.api_key = settings.TON_API_KEY
|
||||
self.timeout = 30
|
||||
|
||||
# HTTP client for API requests
|
||||
self.client = httpx.AsyncClient(
|
||||
timeout=self.timeout,
|
||||
headers={
|
||||
"Authorization": f"Bearer {self.api_key}" if self.api_key else None,
|
||||
"Content-Type": "application/json"
|
||||
}
|
||||
)
|
||||
await send_status("ton_daemon", "working: topup highload wallet")
|
||||
await asyncio.sleep(15)
|
||||
return await main_fn(memory)
|
||||
|
||||
self.cache_manager = get_cache_manager()
|
||||
|
||||
if not highload_state.get('code'):
|
||||
make_log("TON", "Highload wallet contract is not deployed, send deploy transaction..", level="info")
|
||||
created_at_ts = int(datetime.utcnow().timestamp()) - 60
|
||||
await toncenter.send_boc(
|
||||
highload_wallet.create_transfer_message(
|
||||
service_wallet.address.to_string(1, 1, 1),
|
||||
1, HighloadQueryId.from_seqno(0), created_at_ts, send_mode=1, payload="hello world", need_deploy=True
|
||||
)['message'].to_boc(False)
|
||||
)
|
||||
await send_status("ton_daemon", "working: deploying highload wallet")
|
||||
await asyncio.sleep(15)
|
||||
return await main_fn(memory)
|
||||
|
||||
while True:
|
||||
async def close(self):
|
||||
"""Close HTTP client and cleanup resources."""
|
||||
if self.client:
|
||||
await self.client.aclose()
|
||||
|
||||
async def create_wallet(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Create new TON wallet with mnemonic generation.
|
||||
|
||||
Returns:
|
||||
Dict: Wallet creation result with address and private key
|
||||
"""
|
||||
try:
|
||||
sw_seqno_value = await get_sw_seqno()
|
||||
make_log("TON", f"Service running ({sw_seqno_value})", level="debug")
|
||||
|
||||
with db_session() as session:
|
||||
# Проверка отправленных сообщений
|
||||
await send_status("ton_daemon", f"working: processing in-txs (seqno={sw_seqno_value})")
|
||||
async def process_incoming_transaction(transaction: dict):
|
||||
transaction_hash = transaction['transaction_id']['hash']
|
||||
transaction_lt = str(transaction['transaction_id']['lt'])
|
||||
# transaction_success = bool(transaction['success'])
|
||||
|
||||
async def process_incoming_message(blockchain_message: dict):
|
||||
in_msg_cell = Cell.one_from_boc(b64decode(blockchain_message['msg_data']['body']))
|
||||
in_msg_slice = in_msg_cell.refs[0].begin_parse()
|
||||
in_msg_slice.read_uint(32)
|
||||
in_msg_slice.read_uint(8)
|
||||
in_msg_query_id = in_msg_slice.read_uint(23)
|
||||
in_msg_created_at = in_msg_slice.read_uint(64)
|
||||
in_msg_epoch = int(in_msg_created_at // (60 * 60))
|
||||
in_msg_seqno = HighloadQueryId.from_query_id(in_msg_query_id).to_seqno()
|
||||
in_msg_blockchain_task = (
|
||||
session.query(BlockchainTask).filter(
|
||||
and_(
|
||||
BlockchainTask.seqno == in_msg_seqno,
|
||||
BlockchainTask.epoch == in_msg_epoch,
|
||||
)
|
||||
)
|
||||
).first()
|
||||
if not in_msg_blockchain_task:
|
||||
return
|
||||
|
||||
if not (in_msg_blockchain_task.status in ['done']) or in_msg_blockchain_task.transaction_hash != transaction_hash:
|
||||
in_msg_blockchain_task.status = 'done'
|
||||
in_msg_blockchain_task.transaction_hash = transaction_hash
|
||||
in_msg_blockchain_task.transaction_lt = transaction_lt
|
||||
session.commit()
|
||||
|
||||
for blockchain_message in [transaction['in_msg']]:
|
||||
try:
|
||||
await process_incoming_message(blockchain_message)
|
||||
except BaseException as e:
|
||||
pass # make_log("TON_Daemon", f"Error while processing incoming message: {e}" + '\n' + traceback.format_exc(), level='debug')
|
||||
|
||||
# Generate mnemonic phrase
|
||||
mnemonic_response = await self.client.post(
|
||||
f"{self.api_endpoint}/wallet/generate",
|
||||
json={"testnet": self.testnet}
|
||||
)
|
||||
|
||||
if mnemonic_response.status_code != 200:
|
||||
error_msg = f"Failed to generate wallet: {mnemonic_response.text}"
|
||||
await logger.aerror("Wallet generation failed", error=error_msg)
|
||||
return {"error": error_msg}
|
||||
|
||||
mnemonic_data = mnemonic_response.json()
|
||||
|
||||
# Create wallet from mnemonic
|
||||
wallet_response = await self.client.post(
|
||||
f"{self.api_endpoint}/wallet/create",
|
||||
json={
|
||||
"mnemonic": mnemonic_data["mnemonic"],
|
||||
"testnet": self.testnet
|
||||
}
|
||||
)
|
||||
|
||||
if wallet_response.status_code != 200:
|
||||
error_msg = f"Failed to create wallet: {wallet_response.text}"
|
||||
await logger.aerror("Wallet creation failed", error=error_msg)
|
||||
return {"error": error_msg}
|
||||
|
||||
wallet_data = wallet_response.json()
|
||||
|
||||
await logger.ainfo(
|
||||
"Wallet created successfully",
|
||||
address=wallet_data.get("address"),
|
||||
testnet=self.testnet
|
||||
)
|
||||
|
||||
return {
|
||||
"address": wallet_data["address"],
|
||||
"private_key": wallet_data["private_key"],
|
||||
"mnemonic": mnemonic_data["mnemonic"],
|
||||
"testnet": self.testnet
|
||||
}
|
||||
|
||||
except httpx.TimeoutException:
|
||||
error_msg = "Wallet creation timeout"
|
||||
await logger.aerror(error_msg)
|
||||
return {"error": error_msg}
|
||||
except Exception as e:
|
||||
error_msg = f"Wallet creation error: {str(e)}"
|
||||
await logger.aerror("Wallet creation exception", error=str(e))
|
||||
return {"error": error_msg}
|
||||
|
||||
async def get_wallet_balance(self, address: str) -> Dict[str, Any]:
|
||||
"""
|
||||
Get wallet balance with caching for performance.
|
||||
|
||||
Args:
|
||||
address: TON wallet address
|
||||
|
||||
Returns:
|
||||
Dict: Balance information
|
||||
"""
|
||||
try:
|
||||
# Check cache first
|
||||
cache_key = f"ton_balance:{address}"
|
||||
cached_balance = await self.cache_manager.get(cache_key)
|
||||
|
||||
if cached_balance:
|
||||
return cached_balance
|
||||
|
||||
# Fetch from blockchain
|
||||
balance_response = await self.client.get(
|
||||
f"{self.api_endpoint}/wallet/{address}/balance"
|
||||
)
|
||||
|
||||
if balance_response.status_code != 200:
|
||||
error_msg = f"Failed to get balance: {balance_response.text}"
|
||||
return {"error": error_msg}
|
||||
|
||||
balance_data = balance_response.json()
|
||||
|
||||
result = {
|
||||
"balance": int(balance_data.get("balance", 0)), # nanotons
|
||||
"last_transaction_lt": balance_data.get("last_transaction_lt"),
|
||||
"account_state": balance_data.get("account_state", "unknown"),
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
# Cache for 30 seconds
|
||||
await self.cache_manager.set(cache_key, result, ttl=30)
|
||||
|
||||
return result
|
||||
|
||||
except httpx.TimeoutException:
|
||||
return {"error": "Balance fetch timeout"}
|
||||
except Exception as e:
|
||||
await logger.aerror("Balance fetch error", address=address, error=str(e))
|
||||
return {"error": f"Balance fetch error: {str(e)}"}
|
||||
|
||||
async def get_wallet_transactions(
|
||||
self,
|
||||
address: str,
|
||||
limit: int = 20,
|
||||
offset: int = 0
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Get wallet transaction history with pagination.
|
||||
|
||||
Args:
|
||||
address: TON wallet address
|
||||
limit: Number of transactions to fetch
|
||||
offset: Pagination offset
|
||||
|
||||
Returns:
|
||||
Dict: Transaction history
|
||||
"""
|
||||
try:
|
||||
# Check cache
|
||||
cache_key = f"ton_transactions:{address}:{limit}:{offset}"
|
||||
cached_transactions = await self.cache_manager.get(cache_key)
|
||||
|
||||
if cached_transactions:
|
||||
return cached_transactions
|
||||
|
||||
transactions_response = await self.client.get(
|
||||
f"{self.api_endpoint}/wallet/{address}/transactions",
|
||||
params={"limit": limit, "offset": offset}
|
||||
)
|
||||
|
||||
if transactions_response.status_code != 200:
|
||||
error_msg = f"Failed to get transactions: {transactions_response.text}"
|
||||
return {"error": error_msg}
|
||||
|
||||
transactions_data = transactions_response.json()
|
||||
|
||||
result = {
|
||||
"transactions": transactions_data.get("transactions", []),
|
||||
"total": transactions_data.get("total", 0),
|
||||
"limit": limit,
|
||||
"offset": offset,
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
# Cache for 1 minute
|
||||
await self.cache_manager.set(cache_key, result, ttl=60)
|
||||
|
||||
return result
|
||||
|
||||
except httpx.TimeoutException:
|
||||
return {"error": "Transaction fetch timeout"}
|
||||
except Exception as e:
|
||||
await logger.aerror(
|
||||
"Transaction fetch error",
|
||||
address=address,
|
||||
error=str(e)
|
||||
)
|
||||
return {"error": f"Transaction fetch error: {str(e)}"}
|
||||
|
||||
async def send_transaction(
|
||||
self,
|
||||
private_key: str,
|
||||
recipient_address: str,
|
||||
amount: int,
|
||||
message: str = "",
|
||||
**kwargs
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Send TON transaction with validation and monitoring.
|
||||
|
||||
Args:
|
||||
private_key: Encrypted private key
|
||||
recipient_address: Recipient wallet address
|
||||
amount: Amount in nanotons
|
||||
message: Optional message
|
||||
**kwargs: Additional transaction parameters
|
||||
|
||||
Returns:
|
||||
Dict: Transaction result
|
||||
"""
|
||||
try:
|
||||
# Validate inputs
|
||||
if amount <= 0:
|
||||
return {"error": "Amount must be positive"}
|
||||
|
||||
if len(recipient_address) != 48:
|
||||
return {"error": "Invalid recipient address format"}
|
||||
|
||||
# Decrypt private key
|
||||
try:
|
||||
decrypted_key = decrypt_data(private_key, context="wallet")
|
||||
if isinstance(decrypted_key, bytes):
|
||||
decrypted_key = decrypted_key.decode('utf-8')
|
||||
except Exception as e:
|
||||
await logger.aerror("Private key decryption failed", error=str(e))
|
||||
return {"error": "Invalid private key"}
|
||||
|
||||
# Prepare transaction
|
||||
transaction_data = {
|
||||
"private_key": decrypted_key,
|
||||
"recipient": recipient_address,
|
||||
"amount": str(amount),
|
||||
"message": message,
|
||||
"testnet": self.testnet
|
||||
}
|
||||
|
||||
# Send transaction
|
||||
tx_response = await self.client.post(
|
||||
f"{self.api_endpoint}/transaction/send",
|
||||
json=transaction_data
|
||||
)
|
||||
|
||||
if tx_response.status_code != 200:
|
||||
error_msg = f"Transaction failed: {tx_response.text}"
|
||||
await logger.aerror(
|
||||
"Transaction submission failed",
|
||||
recipient=recipient_address,
|
||||
amount=amount,
|
||||
error=error_msg
|
||||
)
|
||||
return {"error": error_msg}
|
||||
|
||||
tx_data = tx_response.json()
|
||||
|
||||
result = {
|
||||
"hash": tx_data["hash"],
|
||||
"lt": tx_data.get("lt"),
|
||||
"fee": tx_data.get("fee", 0),
|
||||
"block_hash": tx_data.get("block_hash"),
|
||||
"timestamp": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
await logger.ainfo(
|
||||
"Transaction sent successfully",
|
||||
hash=result["hash"],
|
||||
recipient=recipient_address,
|
||||
amount=amount
|
||||
)
|
||||
|
||||
return result
|
||||
|
||||
except httpx.TimeoutException:
|
||||
return {"error": "Transaction timeout"}
|
||||
except Exception as e:
|
||||
await logger.aerror(
|
||||
"Transaction send error",
|
||||
recipient=recipient_address,
|
||||
amount=amount,
|
||||
error=str(e)
|
||||
)
|
||||
return {"error": f"Transaction error: {str(e)}"}
|
||||
|
||||
async def get_transaction_status(self, tx_hash: str) -> Dict[str, Any]:
|
||||
"""
|
||||
Get transaction status and confirmation details.
|
||||
|
||||
Args:
|
||||
tx_hash: Transaction hash
|
||||
|
||||
Returns:
|
||||
Dict: Transaction status information
|
||||
"""
|
||||
try:
|
||||
# Check cache
|
||||
cache_key = f"ton_tx_status:{tx_hash}"
|
||||
cached_status = await self.cache_manager.get(cache_key)
|
||||
|
||||
if cached_status and cached_status.get("confirmed"):
|
||||
return cached_status
|
||||
|
||||
status_response = await self.client.get(
|
||||
f"{self.api_endpoint}/transaction/{tx_hash}/status"
|
||||
)
|
||||
|
||||
if status_response.status_code != 200:
|
||||
return {"error": f"Failed to get status: {status_response.text}"}
|
||||
|
||||
status_data = status_response.json()
|
||||
|
||||
result = {
|
||||
"hash": tx_hash,
|
||||
"confirmed": status_data.get("confirmed", False),
|
||||
"failed": status_data.get("failed", False),
|
||||
"confirmations": status_data.get("confirmations", 0),
|
||||
"block_hash": status_data.get("block_hash"),
|
||||
"block_time": status_data.get("block_time"),
|
||||
"fee": status_data.get("fee"),
|
||||
"confirmed_at": status_data.get("confirmed_at"),
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
# Cache confirmed/failed transactions longer
|
||||
cache_ttl = 3600 if result["confirmed"] or result["failed"] else 30
|
||||
await self.cache_manager.set(cache_key, result, ttl=cache_ttl)
|
||||
|
||||
return result
|
||||
|
||||
except httpx.TimeoutException:
|
||||
return {"error": "Status check timeout"}
|
||||
except Exception as e:
|
||||
await logger.aerror("Status check error", tx_hash=tx_hash, error=str(e))
|
||||
return {"error": f"Status check error: {str(e)}"}
|
||||
|
||||
async def validate_address(self, address: str) -> Dict[str, Any]:
|
||||
"""
|
||||
Validate TON address format and existence.
|
||||
|
||||
Args:
|
||||
address: TON address to validate
|
||||
|
||||
Returns:
|
||||
Dict: Validation result
|
||||
"""
|
||||
try:
|
||||
# Basic format validation
|
||||
if len(address) != 48:
|
||||
return {"valid": False, "error": "Invalid address length"}
|
||||
|
||||
# Check against blockchain
|
||||
validation_response = await self.client.post(
|
||||
f"{self.api_endpoint}/address/validate",
|
||||
json={"address": address}
|
||||
)
|
||||
|
||||
if validation_response.status_code != 200:
|
||||
return {"valid": False, "error": "Validation service error"}
|
||||
|
||||
validation_data = validation_response.json()
|
||||
|
||||
return {
|
||||
"valid": validation_data.get("valid", False),
|
||||
"exists": validation_data.get("exists", False),
|
||||
"account_type": validation_data.get("account_type"),
|
||||
"error": validation_data.get("error")
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
await logger.aerror("Address validation error", address=address, error=str(e))
|
||||
return {"valid": False, "error": f"Validation error: {str(e)}"}
|
||||
|
||||
async def get_network_info(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Get TON network information and statistics.
|
||||
|
||||
Returns:
|
||||
Dict: Network information
|
||||
"""
|
||||
try:
|
||||
cache_key = "ton_network_info"
|
||||
cached_info = await self.cache_manager.get(cache_key)
|
||||
|
||||
if cached_info:
|
||||
return cached_info
|
||||
|
||||
network_response = await self.client.get(
|
||||
f"{self.api_endpoint}/network/info"
|
||||
)
|
||||
|
||||
if network_response.status_code != 200:
|
||||
return {"error": f"Failed to get network info: {network_response.text}"}
|
||||
|
||||
network_data = network_response.json()
|
||||
|
||||
result = {
|
||||
"network": "testnet" if self.testnet else "mainnet",
|
||||
"last_block": network_data.get("last_block"),
|
||||
"last_block_time": network_data.get("last_block_time"),
|
||||
"total_accounts": network_data.get("total_accounts"),
|
||||
"total_transactions": network_data.get("total_transactions"),
|
||||
"tps": network_data.get("tps"), # Transactions per second
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
# Cache for 5 minutes
|
||||
await self.cache_manager.set(cache_key, result, ttl=300)
|
||||
|
||||
return result
|
||||
|
||||
except Exception as e:
|
||||
await logger.aerror("Network info error", error=str(e))
|
||||
return {"error": f"Network info error: {str(e)}"}
|
||||
|
||||
async def estimate_transaction_fee(
|
||||
self,
|
||||
sender_address: str,
|
||||
recipient_address: str,
|
||||
amount: int,
|
||||
message: str = ""
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Estimate transaction fee before sending.
|
||||
|
||||
Args:
|
||||
sender_address: Sender wallet address
|
||||
recipient_address: Recipient wallet address
|
||||
amount: Amount in nanotons
|
||||
message: Optional message
|
||||
|
||||
Returns:
|
||||
Dict: Fee estimation
|
||||
"""
|
||||
try:
|
||||
fee_response = await self.client.post(
|
||||
f"{self.api_endpoint}/transaction/estimate-fee",
|
||||
json={
|
||||
"sender": sender_address,
|
||||
"recipient": recipient_address,
|
||||
"amount": str(amount),
|
||||
"message": message
|
||||
}
|
||||
)
|
||||
|
||||
if fee_response.status_code != 200:
|
||||
return {"error": f"Fee estimation failed: {fee_response.text}"}
|
||||
|
||||
fee_data = fee_response.json()
|
||||
|
||||
return {
|
||||
"estimated_fee": fee_data.get("fee", 0),
|
||||
"estimated_fee_tons": str(Decimal(fee_data.get("fee", 0)) / Decimal("1000000000")),
|
||||
"gas_used": fee_data.get("gas_used"),
|
||||
"message_size": len(message.encode('utf-8')),
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
await logger.aerror("Fee estimation error", error=str(e))
|
||||
return {"error": f"Fee estimation error: {str(e)}"}
|
||||
|
||||
async def monitor_transaction(self, tx_hash: str, max_wait_time: int = 300) -> Dict[str, Any]:
|
||||
"""
|
||||
Monitor transaction until confirmation or timeout.
|
||||
|
||||
Args:
|
||||
tx_hash: Transaction hash to monitor
|
||||
max_wait_time: Maximum wait time in seconds
|
||||
|
||||
Returns:
|
||||
Dict: Final transaction status
|
||||
"""
|
||||
start_time = datetime.utcnow()
|
||||
check_interval = 5 # Check every 5 seconds
|
||||
|
||||
while (datetime.utcnow() - start_time).seconds < max_wait_time:
|
||||
status = await self.get_transaction_status(tx_hash)
|
||||
|
||||
if status.get("error"):
|
||||
return status
|
||||
|
||||
if status.get("confirmed") or status.get("failed"):
|
||||
await logger.ainfo(
|
||||
"Transaction monitoring completed",
|
||||
tx_hash=tx_hash,
|
||||
confirmed=status.get("confirmed"),
|
||||
failed=status.get("failed"),
|
||||
duration=(datetime.utcnow() - start_time).seconds
|
||||
)
|
||||
return status
|
||||
|
||||
await asyncio.sleep(check_interval)
|
||||
|
||||
# Timeout reached
|
||||
await logger.awarning(
|
||||
"Transaction monitoring timeout",
|
||||
tx_hash=tx_hash,
|
||||
max_wait_time=max_wait_time
|
||||
)
|
||||
|
||||
return {
|
||||
"hash": tx_hash,
|
||||
"confirmed": False,
|
||||
"timeout": True,
|
||||
"error": "Monitoring timeout reached"
|
||||
}
|
||||
|
||||
async def get_smart_contract_info(self, address: str) -> Dict[str, Any]:
|
||||
"""
|
||||
Get smart contract information and ABI.
|
||||
|
||||
Args:
|
||||
address: Smart contract address
|
||||
|
||||
Returns:
|
||||
Dict: Contract information
|
||||
"""
|
||||
try:
|
||||
cache_key = f"ton_contract:{address}"
|
||||
cached_info = await self.cache_manager.get(cache_key)
|
||||
|
||||
if cached_info:
|
||||
return cached_info
|
||||
|
||||
contract_response = await self.client.get(
|
||||
f"{self.api_endpoint}/contract/{address}/info"
|
||||
)
|
||||
|
||||
if contract_response.status_code != 200:
|
||||
return {"error": f"Failed to get contract info: {contract_response.text}"}
|
||||
|
||||
contract_data = contract_response.json()
|
||||
|
||||
result = {
|
||||
"address": address,
|
||||
"contract_type": contract_data.get("contract_type"),
|
||||
"is_verified": contract_data.get("is_verified", False),
|
||||
"abi": contract_data.get("abi"),
|
||||
"source_code": contract_data.get("source_code"),
|
||||
"compiler_version": contract_data.get("compiler_version"),
|
||||
"deployment_block": contract_data.get("deployment_block"),
|
||||
"updated_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
# Cache for 1 hour
|
||||
await self.cache_manager.set(cache_key, result, ttl=3600)
|
||||
|
||||
return result
|
||||
|
||||
except Exception as e:
|
||||
await logger.aerror("Contract info error", address=address, error=str(e))
|
||||
return {"error": f"Contract info error: {str(e)}"}
|
||||
|
||||
async def call_smart_contract(
|
||||
self,
|
||||
contract_address: str,
|
||||
method: str,
|
||||
params: Dict[str, Any],
|
||||
private_key: Optional[str] = None
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Call smart contract method.
|
||||
|
||||
Args:
|
||||
contract_address: Contract address
|
||||
method: Method name to call
|
||||
params: Method parameters
|
||||
private_key: Private key for write operations
|
||||
|
||||
Returns:
|
||||
Dict: Contract call result
|
||||
"""
|
||||
try:
|
||||
call_data = {
|
||||
"contract": contract_address,
|
||||
"method": method,
|
||||
"params": params
|
||||
}
|
||||
|
||||
# Add private key for write operations
|
||||
if private_key:
|
||||
try:
|
||||
sw_transactions = await toncenter.get_transactions(highload_wallet.address.to_string(1, 1, 1), limit=100)
|
||||
for sw_transaction in sw_transactions:
|
||||
try:
|
||||
await process_incoming_transaction(sw_transaction)
|
||||
except BaseException as e:
|
||||
make_log("TON_Daemon", f"Error while processing incoming transaction: {e}", level="debug")
|
||||
except BaseException as e:
|
||||
make_log("TON_Daemon", f"Error while getting service wallet transactions: {e}", level="ERROR")
|
||||
|
||||
await send_status("ton_daemon", f"working: processing out-txs (seqno={sw_seqno_value})")
|
||||
# Отправка подписанных сообщений
|
||||
for blockchain_task in (
|
||||
session.query(BlockchainTask).filter(
|
||||
BlockchainTask.status == 'processing',
|
||||
).order_by(BlockchainTask.updated.asc()).all()
|
||||
):
|
||||
make_log("TON_Daemon", f"Processing task (processing) {blockchain_task.id}")
|
||||
query_boc = bytes.fromhex(blockchain_task.meta['signed_message'])
|
||||
errors_list = []
|
||||
|
||||
try:
|
||||
await toncenter.send_boc(query_boc)
|
||||
except BaseException as e:
|
||||
errors_list.append(f"{e}")
|
||||
|
||||
try:
|
||||
make_log("TON_Daemon", str(
|
||||
httpx.post(
|
||||
'https://tonapi.io/v2/blockchain/message',
|
||||
json={
|
||||
'boc': query_boc.hex()
|
||||
}
|
||||
).text
|
||||
))
|
||||
except BaseException as e:
|
||||
make_log("TON_Daemon", f"Error while pushing task to tonkeeper ({blockchain_task.id}): {e}", level="ERROR")
|
||||
errors_list.append(f"{e}")
|
||||
|
||||
blockchain_task.updated = datetime.utcnow()
|
||||
|
||||
if blockchain_task.meta['sign_created'] + 10 * 60 < datetime.utcnow().timestamp():
|
||||
# or sum([int("terminating vm with exit code 36" in e) for e in errors_list]) > 0:
|
||||
make_log("TON_Daemon", f"Task {blockchain_task.id} done", level="DEBUG")
|
||||
blockchain_task.status = 'done'
|
||||
session.commit()
|
||||
continue
|
||||
|
||||
await asyncio.sleep(0.5)
|
||||
|
||||
await send_status("ton_daemon", f"working: creating new messages (seqno={sw_seqno_value})")
|
||||
# Создание новых подписей
|
||||
for blockchain_task in (
|
||||
session.query(BlockchainTask).filter(BlockchainTask.status == 'wait').all()
|
||||
):
|
||||
try:
|
||||
# Check processing tasks in current epoch < 3_000_000
|
||||
if (
|
||||
session.query(BlockchainTask).filter(
|
||||
BlockchainTask.epoch == blockchain_task.epoch,
|
||||
).count() > 3_000_000
|
||||
):
|
||||
make_log("TON", f"Too many processing tasks in epoch {blockchain_task.epoch}", level="error")
|
||||
await send_status("ton_daemon", f"working: too many tasks in epoch {blockchain_task.epoch}")
|
||||
await asyncio.sleep(5)
|
||||
continue
|
||||
|
||||
sign_created = int(datetime.utcnow().timestamp()) - 60
|
||||
try:
|
||||
current_epoch = int(datetime.utcnow().timestamp() // (60 * 60))
|
||||
max_epoch_seqno = (
|
||||
session.query(func.max(BlockchainTask.seqno)).filter(
|
||||
BlockchainTask.epoch == current_epoch
|
||||
).scalar() or 0
|
||||
)
|
||||
current_epoch_shift = 3_000_000 if current_epoch % 2 == 0 else 0
|
||||
current_seqno = max_epoch_seqno + 1 + (current_epoch_shift if max_epoch_seqno == 0 else 0)
|
||||
except BaseException as e:
|
||||
make_log("CRITICAL", f"Error calculating epoch,seqno: {e}", level="error")
|
||||
current_epoch = 0
|
||||
current_seqno = 0
|
||||
|
||||
blockchain_task.seqno = current_seqno
|
||||
blockchain_task.epoch = current_epoch
|
||||
blockchain_task.status = 'processing'
|
||||
try:
|
||||
query = highload_wallet.create_transfer_message(
|
||||
blockchain_task.destination, int(blockchain_task.amount), HighloadQueryId.from_seqno(current_seqno),
|
||||
sign_created, send_mode=1,
|
||||
payload=Cell.one_from_boc(b64decode(blockchain_task.payload))
|
||||
)
|
||||
query_boc = query['message'].to_boc(False)
|
||||
except BaseException as e:
|
||||
make_log("TON", f"Error creating transfer message: {e}", level="error")
|
||||
query_boc = begin_cell().end_cell().to_boc(False)
|
||||
|
||||
blockchain_task.meta = {
|
||||
**blockchain_task.meta,
|
||||
'sign_created': sign_created,
|
||||
'signed_message': query_boc.hex(),
|
||||
}
|
||||
session.commit()
|
||||
make_log("TON", f"Created signed message for task {blockchain_task.id}" + '\n' + traceback.format_exc(), level="info")
|
||||
except BaseException as e:
|
||||
make_log("TON", f"Error processing task {blockchain_task.id}: {e}" + '\n' + traceback.format_exc(), level="error")
|
||||
continue
|
||||
|
||||
await asyncio.sleep(1)
|
||||
|
||||
await asyncio.sleep(1)
|
||||
await send_status("ton_daemon", f"working (seqno={sw_seqno_value})")
|
||||
except BaseException as e:
|
||||
make_log("TON", f"Error: {e}", level="error")
|
||||
await asyncio.sleep(3)
|
||||
|
||||
# if __name__ == '__main__':
|
||||
# loop = asyncio.get_event_loop()
|
||||
# loop.run_until_complete(main())
|
||||
# loop.close()
|
||||
|
||||
decrypted_key = decrypt_data(private_key, context="wallet")
|
||||
if isinstance(decrypted_key, bytes):
|
||||
decrypted_key = decrypted_key.decode('utf-8')
|
||||
call_data["private_key"] = decrypted_key
|
||||
except Exception as e:
|
||||
return {"error": "Invalid private key"}
|
||||
|
||||
contract_response = await self.client.post(
|
||||
f"{self.api_endpoint}/contract/call",
|
||||
json=call_data
|
||||
)
|
||||
|
||||
if contract_response.status_code != 200:
|
||||
return {"error": f"Contract call failed: {contract_response.text}"}
|
||||
|
||||
call_result = contract_response.json()
|
||||
|
||||
await logger.ainfo(
|
||||
"Smart contract called",
|
||||
contract=contract_address,
|
||||
method=method,
|
||||
success=call_result.get("success", False)
|
||||
)
|
||||
|
||||
return call_result
|
||||
|
||||
except Exception as e:
|
||||
await logger.aerror(
|
||||
"Contract call error",
|
||||
contract=contract_address,
|
||||
method=method,
|
||||
error=str(e)
|
||||
)
|
||||
return {"error": f"Contract call error: {str(e)}"}
|
||||
|
||||
# Global TON service instance
|
||||
_ton_service = None
|
||||
|
||||
async def get_ton_service() -> TONService:
|
||||
"""Get or create global TON service instance."""
|
||||
global _ton_service
|
||||
if _ton_service is None:
|
||||
_ton_service = TONService()
|
||||
return _ton_service
|
||||
|
||||
async def cleanup_ton_service():
|
||||
"""Cleanup global TON service instance."""
|
||||
global _ton_service
|
||||
if _ton_service:
|
||||
await _ton_service.close()
|
||||
_ton_service = None
|
||||
Reference in new issue
Block a user