23 Commits
Author SHA1 Message Date
user 14fa6fe096 q 2025-05-30 18:56:53 +03:00
user 4605a46765 nice text 2025-05-29 21:01:24 +03:00
user 6de350c2a3 add qq command 2025-05-04 19:14:13 +03:00
user 11e2645aa9 edit player ui 2025-05-03 16:00:30 +03:00
user 3d9e7c7966 fix 2025-04-26 13:29:31 +03:00
user 4483c194da fix magic 2025-04-26 13:24:35 +03:00
user d2ff332490 fix convert service 2025-04-26 13:13:47 +03:00
user 2476d3b3b6 fix someth 2025-04-26 13:01:10 +03:00
user 0e4268fb4d python magic for images 2025-04-26 12:52:32 +03:00
user 0586ed9d94 disable response header credentials 2025-03-19 14:29:32 +03:00
user a6750fe35c fix misprint 2025-03-18 14:40:54 +03:00
user 3ab014d358 fix selectWallet 2025-03-17 13:54:36 +03:00
user 4a739e4b1b fix misprint 2025-03-17 13:47:22 +03:00
user 1a1f4301cf downloadble fix 2025-03-17 13:43:31 +03:00
user b02ded982c dowload & selectWallet & fixes & auth me 2025-03-17 13:40:37 +03:00
user ac6d102a3f downloadable, selectWallet 2025-03-14 13:57:30 +03:00
user 66da3541e7 try fix licenses 2025-03-14 10:54:57 +03:00
user 7d962d463b add free upload hint 2025-03-14 00:41:31 +03:00
user 4643d7f202 fix misprint 2025-03-13 21:00:04 +03:00
user 54bc545090 add status for ton_daemon 2025-03-13 20:14:40 +03:00
user b63f663bd2 set newContent amount task 2025-03-13 20:01:48 +03:00
user a35481fa71 fix misprint 2025-03-13 19:47:04 +03:00
user 276a09fbf6 restart policy 2025-03-13 18:41:55 +03:00
22 changed files with 226 additions and 112 deletions

No files matched your search

+1 -1
View File
@@ -21,7 +21,7 @@ RUN apt-get update && apt-get install -y \
apt-get update && \ apt-get update && \
apt-get install -y docker-ce-cli apt-get install -y docker-ce-cli
RUN apt-get install -y ffmpeg RUN apt-get install libmagic1 -y
CMD ["python", "app"] CMD ["python", "app"]
+3 -1
View File
@@ -13,7 +13,7 @@ app.register_middleware(close_db_session, "response")
from app.api.routes._index import s_index, s_favicon from app.api.routes._index import s_index, s_favicon
from app.api.routes._system import s_api_v1_node, s_api_system_version, s_api_system_send_status, s_api_v1_node_friendly from app.api.routes._system import s_api_v1_node, s_api_system_version, s_api_system_send_status, s_api_v1_node_friendly
from app.api.routes.auth import s_api_v1_auth_twa from app.api.routes.auth import s_api_v1_auth_twa, s_api_v1_auth_select_wallet, s_api_v1_auth_me
from app.api.routes.statics import s_api_tonconnect_manifest, s_api_platform_metadata from app.api.routes.statics import s_api_tonconnect_manifest, s_api_platform_metadata
from app.api.routes.node_storage import s_api_v1_storage_post, s_api_v1_storage_get, \ from app.api.routes.node_storage import s_api_v1_storage_post, s_api_v1_storage_get, \
s_api_v1_storage_decode_cid s_api_v1_storage_decode_cid
@@ -37,6 +37,8 @@ app.add_route(s_api_tonconnect_manifest, "/api/tonconnect-manifest.json", method
app.add_route(s_api_platform_metadata, "/api/platform-metadata.json", methods=["GET", "OPTIONS"]) app.add_route(s_api_platform_metadata, "/api/platform-metadata.json", methods=["GET", "OPTIONS"])
app.add_route(s_api_v1_auth_twa, "/api/v1/auth.twa", methods=["POST", "OPTIONS"]) app.add_route(s_api_v1_auth_twa, "/api/v1/auth.twa", methods=["POST", "OPTIONS"])
app.add_route(s_api_v1_auth_me, "/api/v1/auth.me", methods=["GET", "OPTIONS"])
app.add_route(s_api_v1_auth_select_wallet, "/api/v1/auth.selectWallet", methods=["POST", "OPTIONS"])
app.add_route(s_api_v1_tonconnect_new, "/api/v1/tonconnect.new", methods=["GET", "OPTIONS"]) app.add_route(s_api_v1_tonconnect_new, "/api/v1/tonconnect.new", methods=["GET", "OPTIONS"])
app.add_route(s_api_v1_tonconnect_logout, "/api/v1/tonconnect.logout", methods=["POST", "OPTIONS"]) app.add_route(s_api_v1_tonconnect_logout, "/api/v1/tonconnect.logout", methods=["POST", "OPTIONS"])
+1 -1
View File
@@ -16,7 +16,7 @@ def attach_headers(response):
response.headers["Access-Control-Allow-Origin"] = "*" response.headers["Access-Control-Allow-Origin"] = "*"
response.headers["Access-Control-Allow-Methods"] = "GET, POST, OPTIONS" response.headers["Access-Control-Allow-Methods"] = "GET, POST, OPTIONS"
response.headers["Access-Control-Allow-Headers"] = "Origin, Content-Type, Accept, Authorization, Referer, User-Agent, Sec-Fetch-Dest, Sec-Fetch-Mode, Sec-Fetch-Site, x-file-name, x-last-chunk, x-chunk-start, x-upload-id" response.headers["Access-Control-Allow-Headers"] = "Origin, Content-Type, Accept, Authorization, Referer, User-Agent, Sec-Fetch-Dest, Sec-Fetch-Mode, Sec-Fetch-Site, x-file-name, x-last-chunk, x-chunk-start, x-upload-id"
response.headers["Access-Control-Allow-Credentials"] = "true" # response.headers["Access-Control-Allow-Credentials"] = "true"
return response return response
+9 -2
View File
@@ -90,7 +90,8 @@ async def s_api_v1_blockchain_send_new_content_message(request):
title=content_title, title=content_title,
cover_url=f"{PROJECT_HOST}/api/v1.5/storage/{image_content_cid.serialize_v2()}" if image_content_cid else None, cover_url=f"{PROJECT_HOST}/api/v1.5/storage/{image_content_cid.serialize_v2()}" if image_content_cid else None,
authors=request.json['authors'], authors=request.json['authors'],
hashtags=request.json['hashtags'] hashtags=request.json['hashtags'],
downloadable=request.json['downloadable'] if 'downloadable' in request.json else False,
) )
royalties_dict = begin_dict(8) royalties_dict = begin_dict(8)
@@ -133,6 +134,7 @@ async def s_api_v1_blockchain_send_new_content_message(request):
blockchain_task = BlockchainTask( blockchain_task = BlockchainTask(
destination=platform.address.to_string(1, 1, 1), destination=platform.address.to_string(1, 1, 1),
amount=str(int(0.03 * 10 ** 9)),
payload=b64encode( payload=b64encode(
begin_cell() begin_cell()
.store_uint(0x5491d08c, 32) .store_uint(0x5491d08c, 32)
@@ -187,7 +189,9 @@ async def s_api_v1_blockchain_send_new_content_message(request):
} }
) )
return response.json({ return response.json({
'promoUpload': True, 'address': "free",
'amount': str(int(0.03 * 10 ** 9)),
'payload': ""
}) })
await request.ctx.user_uploader_wrapper.send_message( await request.ctx.user_uploader_wrapper.send_message(
@@ -254,6 +258,9 @@ async def s_api_v1_blockchain_send_purchase_content_message(request):
assert field_key in request.json, f"No {field_key} provided" assert field_key in request.json, f"No {field_key} provided"
assert field_value(request.json[field_key]), f"Invalid {field_key} provided" assert field_value(request.json[field_key]), f"Invalid {field_key} provided"
if not request.ctx.user.wallet_address(request.ctx.db_session):
return response.json({"error": "No wallet address provided"}, status=400)
license_exist = request.ctx.db_session.query(UserContent).filter_by( license_exist = request.ctx.db_session.query(UserContent).filter_by(
onchain_address=request.json['content_address'], onchain_address=request.json['content_address'],
).first() ).first()
+1 -1
View File
@@ -54,7 +54,7 @@ async def s_api_v1_node_friendly(request):
for service_key, service in request.app.ctx.memory.known_states.items(): for service_key, service in request.app.ctx.memory.known_states.items():
response_plain_text += f""" response_plain_text += f"""
{service_key}: {service_key}:
status: {service['status'] if (service['timestamp'] and (datetime.now() - service['timestamp']).total_seconds() < 30) else 'not working: timeout'} status: {service['status'] if (service['timestamp'] and (datetime.now() - service['timestamp']).total_seconds() < 120) else 'not working: timeout'}
delay: {round((datetime.now() - service['timestamp']).total_seconds(), 3) if service['timestamp'] else -1} delay: {round((datetime.now() - service['timestamp']).total_seconds(), 3) if service['timestamp'] else -1}
""" """
return response.text(response_plain_text, content_type='text/plain') return response.text(response_plain_text, content_type='text/plain')
+71 -2
View File
@@ -1,4 +1,5 @@
from datetime import datetime from datetime import datetime
from uuid import uuid4
from aiogram.utils.web_app import safe_parse_webapp_init_data from aiogram.utils.web_app import safe_parse_webapp_init_data
from sanic import response from sanic import response
@@ -109,8 +110,7 @@ async def s_api_v1_auth_twa(request):
WalletConnection.network == 'ton', WalletConnection.network == 'ton',
WalletConnection.invalidated == False WalletConnection.invalidated == False
) )
))).scalars().first() ).order_by(WalletConnection.created.desc()))).scalars().first()
known_user.last_use = datetime.now() known_user.last_use = datetime.now()
request.ctx.db_session.commit() request.ctx.db_session.commit()
@@ -119,3 +119,72 @@ async def s_api_v1_auth_twa(request):
'connected_wallet': ton_connection.json_format() if ton_connection else None, 'connected_wallet': ton_connection.json_format() if ton_connection else None,
'auth_v1_token': new_user_key['auth_v1_token'] 'auth_v1_token': new_user_key['auth_v1_token']
}) })
async def s_api_v1_auth_me(request):
if not request.ctx.user:
return response.json({"error": "Unauthorized"}, status=401)
ton_connection = (request.ctx.db_session.execute(
select(WalletConnection).where(
and_(
WalletConnection.user_id == request.ctx.user.id,
WalletConnection.network == 'ton',
WalletConnection.invalidated == False
)
).order_by(WalletConnection.created.desc())
)).scalars().first()
return response.json({
'user': request.ctx.user.json_format(),
'connected_wallet': ton_connection.json_format() if ton_connection else None
})
async def s_api_v1_auth_select_wallet(request):
if not request.ctx.user:
return response.json({"error": "Unauthorized"}, status=401)
try:
data = request.json
except Exception as e:
return response.json({"error": "Invalid JSON"}, status=400)
if "wallet_address" not in data:
return response.json({"error": "wallet_address is required"}, status=400)
# Convert raw wallet address to canonical format using Address from tonsdk.utils
raw_addr = data["wallet_address"]
canonical_address = Address(raw_addr).to_string(1, 1, 1)
db_session = request.ctx.db_session
user = request.ctx.user
# Check if a WalletConnection already exists for this user with the given canonical wallet address
existing_connection = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.wallet_address == canonical_address
).first()
if not existing_connection:
return response.json({"error": "Wallet connection not found"}, status=404)
saved_values = {
'keys': existing_connection.keys,
'meta': existing_connection.meta,
'wallet_key': existing_connection.wallet_key,
'connection_id': existing_connection.connection_id + uuid4().hex,
'network': existing_connection.network,
}
new_connection = WalletConnection(
**saved_values,
user_id=user.id,
wallet_address=canonical_address,
created=datetime.now(),
updated=datetime.now(),
invalidated=False,
without_pk=False
)
db_session.add(new_connection)
db_session.commit()
return response.empty(status=200)
+6 -1
View File
@@ -147,6 +147,7 @@ async def s_api_v1_content_view(request, content_address: str):
).first() ).first()
if converted_content: if converted_content:
display_options['content_url'] = converted_content.web_url display_options['content_url'] = converted_content.web_url
opts['content_ext'] = converted_content.filename.split('.')[-1]
content_meta = content['encrypted_content'].json_format() content_meta = content['encrypted_content'].json_format()
content_metadata = StoredContent.from_cid(request.ctx.db_session, content_meta.get('metadata_cid') or None) content_metadata = StoredContent.from_cid(request.ctx.db_session, content_meta.get('metadata_cid') or None)
@@ -154,11 +155,15 @@ async def s_api_v1_content_view(request, content_address: str):
content_metadata_json = json.loads(f.read()) content_metadata_json = json.loads(f.read())
display_options['metadata'] = content_metadata_json display_options['metadata'] = content_metadata_json
opts['downloadable'] = content_metadata_json.get('downloadable', False)
if opts['downloadable']:
if not ('listen' in opts['have_licenses']):
opts['downloadable'] = False
return response.json({ return response.json({
**opts, **opts,
'encrypted': content['encrypted_content'].json_format(), 'encrypted': content['encrypted_content'].json_format(),
'display_options': display_options 'display_options': display_options,
}) })
+3 -3
View File
@@ -162,7 +162,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
'mov': 'video' 'mov': 'video'
}.get(preview_content.filename.split('.')[-1], content_type_declared) }.get(preview_content.filename.split('.')[-1], content_type_declared)
hashtags_str = (' '.join(f"#{_h}" for _h in metadata_content_json.get('hashtags', []))).strip() hashtags_str = metadata_content_json.get('description', '').strip()
if hashtags_str: if hashtags_str:
hashtags_str = hashtags_str + '\n' hashtags_str = hashtags_str + '\n'
@@ -212,7 +212,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
types.InlineQueryResultCachedAudio( types.InlineQueryResultCachedAudio(
id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}", id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}",
audio_file_id=decrypted_content.meta['telegram_file_cache_preview'], audio_file_id=decrypted_content.meta['telegram_file_cache_preview'],
caption=hashtags_str + user.translated('p_playerContext_preview'), caption=hashtags_str + user.translated('p_playerAudioContext_preview'),
parse_mode='html', parse_mode='html',
reply_markup=get_inline_keyboard([ reply_markup=get_inline_keyboard([
[ [
@@ -242,7 +242,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}", id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}",
video_file_id=decrypted_content.meta['telegram_file_cache_preview'], video_file_id=decrypted_content.meta['telegram_file_cache_preview'],
title=title, title=title,
caption=hashtags_str + user.translated('p_playerContext_preview'), caption=hashtags_str + user.translated('p_playerVideoContext_preview'),
parse_mode='html', parse_mode='html',
reply_markup=get_inline_keyboard([ reply_markup=get_inline_keyboard([
[ [
+2 -2
View File
@@ -1,10 +1,10 @@
from app.core.models import Asset from app.core.models import BlockchainTask
from app.core.models.base import AlchemyBase from app.core.models.base import AlchemyBase
def create_maria_tables(engine): def create_maria_tables(engine):
"""Create all tables in the database.""" """Create all tables in the database."""
Asset() BlockchainTask()
AlchemyBase.metadata.create_all(engine) AlchemyBase.metadata.create_all(engine)
+91 -54
View File
@@ -4,6 +4,7 @@ import os
import uuid import uuid
import json import json
import shutil import shutil
import magic # python-magic for MIME detection
from base58 import b58decode, b58encode from base58 import b58decode, b58encode
from sqlalchemy import and_, or_ from sqlalchemy import and_, or_
from app.core.models.node_storage import StoredContent from app.core.models.node_storage import StoredContent
@@ -33,12 +34,7 @@ async def convert_loop(memory):
make_log("ConvertProcess", "No content to convert", level="debug") make_log("ConvertProcess", "No content to convert", level="debug")
return 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} with preview interval {preview_interval}", level="info")
decrypted_content = session.query(StoredContent).filter( decrypted_content = session.query(StoredContent).filter(
StoredContent.id == unprocessed_encrypted_content.decrypted_content_id StoredContent.id == unprocessed_encrypted_content.decrypted_content_id
).first() ).first()
@@ -46,35 +42,80 @@ async def convert_loop(memory):
make_log("ConvertProcess", "Decrypted content not found", level="error") make_log("ConvertProcess", "Decrypted content not found", level="error")
return return
# Определяем путь и расширение входного файла
# List of conversion options to process
REQUIRED_CONVERT_OPTIONS = ['high', 'low', 'low_preview']
converted_content = {} # Mapping: option -> sha256 hash of output file
# Define input file path and extract its extension from filename
input_file_path = f"/Storage/storedContent/{decrypted_content.hash}" 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 ====
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 = ""
input_ext = unprocessed_encrypted_content.filename.split('.')[-1] if '.' in unprocessed_encrypted_content.filename else "mp4" if mime_type.startswith("video/"):
content_kind = "video"
elif mime_type.startswith("audio/"):
content_kind = "audio"
else:
content_kind = "other"
# Logs directory mapping 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" logs_dir = "/Storage/logs/converter"
# Process each conversion option in sequence
for option in REQUIRED_CONVERT_OPTIONS: for option in REQUIRED_CONVERT_OPTIONS:
# Set quality parameter and trim option (only for preview) # Set quality parameter and trim option (only for preview)
if option == "low_preview": if option == "low_preview":
quality = "low" quality = "low"
trim_value = f"{preview_interval[0]}-{preview_interval[1]}" trim_value = f"{preview_interval[0]}-{preview_interval[1]}"
else: else:
quality = option # 'high' or 'low' quality = option
trim_value = None trim_value = None
# Generate a unique output directory for docker container # Generate a unique output directory for docker container
output_uuid = str(uuid.uuid4()) output_uuid = str(uuid.uuid4())
output_dir = f"/Storage/storedContent/converter-output/{output_uuid}" output_dir = f"/Storage/storedContent/converter-output/{output_uuid}"
# Build the docker command with appropriate volume mounts and parameters # Build the docker command
cmd = [ cmd = [
"docker", "run", "--rm", "docker", "run", "--rm",
"-v", f"{input_file_path}:/app/input", "-v", f"{input_file_path}:/app/input",
@@ -86,8 +127,9 @@ async def convert_loop(memory):
] ]
if trim_value: if trim_value:
cmd.extend(["--trim", trim_value]) cmd.extend(["--trim", trim_value])
if content_kind == "audio":
cmd.append("--audio-only") # audio-only flag
# Run the docker container asynchronously
process = await asyncio.create_subprocess_exec( process = await asyncio.create_subprocess_exec(
*cmd, *cmd,
stdout=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE,
@@ -98,22 +140,24 @@ async def convert_loop(memory):
make_log("ConvertProcess", f"Docker conversion failed for option {option}: {stderr.decode()}", level="error") make_log("ConvertProcess", f"Docker conversion failed for option {option}: {stderr.decode()}", level="error")
return return
# List files in the output directory # List files in output dir
try: try:
files = os.listdir(output_dir.replace("/Storage/storedContent", "/app/data")) files = os.listdir(output_dir.replace("/Storage/storedContent", "/app/data"))
except Exception as e: except Exception as e:
make_log("ConvertProcess", f"Error reading output directory {output_dir}: {e}", level="error") make_log("ConvertProcess", f"Error reading output directory {output_dir}: {e}", level="error")
return return
# Exclude 'output.json' and expect exactly one media output file
media_files = [f for f in files if f != "output.json"] media_files = [f for f in files if f != "output.json"]
if len(media_files) != 1: if len(media_files) != 1:
make_log("ConvertProcess", f"Expected one media file, found {len(media_files)} for option {option}", level="error") make_log("ConvertProcess", f"Expected one media file, found {len(media_files)} for option {option}", level="error")
return return
output_file = os.path.join(output_dir.replace("/Storage/storedContent", "/app/data"), media_files[0]) output_file = os.path.join(
output_dir.replace("/Storage/storedContent", "/app/data"),
media_files[0]
)
# Compute SHA256 hash of the output file using async subprocess # Compute SHA256 hash of the output file
hash_process = await asyncio.create_subprocess_exec( hash_process = await asyncio.create_subprocess_exec(
"sha256sum", output_file, "sha256sum", output_file,
stdout=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE,
@@ -126,6 +170,7 @@ async def convert_loop(memory):
file_hash = hash_stdout.decode().split()[0] file_hash = hash_stdout.decode().split()[0]
file_hash = b58encode(bytes.fromhex(file_hash)).decode() file_hash = b58encode(bytes.fromhex(file_hash)).decode()
# Save new StoredContent if not exists
if not session.query(StoredContent).filter( if not session.query(StoredContent).filter(
StoredContent.hash == file_hash StoredContent.hash == file_hash
).first(): ).first():
@@ -134,9 +179,7 @@ async def convert_loop(memory):
hash=file_hash, hash=file_hash,
user_id=unprocessed_encrypted_content.user_id, user_id=unprocessed_encrypted_content.user_id,
filename=media_files[0], filename=media_files[0],
meta={ meta={'encrypted_file_hash': unprocessed_encrypted_content.hash},
'encrypted_file_hash': unprocessed_encrypted_content.hash,
},
created=datetime.now(), created=datetime.now(),
) )
session.add(new_content) session.add(new_content)
@@ -156,41 +199,32 @@ async def convert_loop(memory):
converted_content[option] = file_hash converted_content[option] = file_hash
# Process output.json: read its contents and update meta['ffprobe_meta'] # Process output.json for ffprobe_meta
output_json_path = os.path.join(output_dir.replace("/Storage/storedContent", "/app/data"), "output.json") output_json_path = os.path.join(
if os.path.exists(output_json_path): output_dir.replace("/Storage/storedContent", "/app/data"),
if unprocessed_encrypted_content.meta.get('ffprobe_meta') is None: "output.json"
)
if os.path.exists(output_json_path) and unprocessed_encrypted_content.meta.get('ffprobe_meta') is None:
try: try:
with open(output_json_path, "r") as f: with open(output_json_path, "r") as f:
output_json_content = f.read() ffprobe_meta = json.load(f)
except Exception as e:
make_log("ConvertProcess", f"Error reading output.json for option {option}: {e}", level="error")
return
try:
ffprobe_meta = json.loads(output_json_content)
except Exception as e:
make_log("ConvertProcess", f"Error parsing output.json for option {option}: {e}", level="error")
return
unprocessed_encrypted_content.meta = { unprocessed_encrypted_content.meta = {
**unprocessed_encrypted_content.meta, **unprocessed_encrypted_content.meta,
'ffprobe_meta': ffprobe_meta 'ffprobe_meta': ffprobe_meta
} }
else: except Exception as e:
make_log("ConvertProcess", f"output.json not found for option {option}", level="error") make_log("ConvertProcess", f"Error handling output.json for option {option}: {e}", level="error")
# Remove the output directory after processing # Cleanup output directory
try: try:
shutil.rmtree(output_dir.replace("/Storage/storedContent", "/app/data")) shutil.rmtree(output_dir.replace("/Storage/storedContent", "/app/data"))
except Exception as e: except Exception as e:
make_log("ConvertProcess", f"Error removing output directory {output_dir}: {e}", level="error") make_log("ConvertProcess", f"Error removing output dir {output_dir}: {e}", level="warning")
# Continue even if deletion fails
# Finalize original record
make_log("ConvertProcess", f"Content {unprocessed_encrypted_content.id} processed. Converted content: {converted_content}", level="info") make_log("ConvertProcess", f"Content {unprocessed_encrypted_content.id} processed. Converted content: {converted_content}", level="info")
unprocessed_encrypted_content.btfs_cid = ContentId( unprocessed_encrypted_content.btfs_cid = ContentId(
version=2, content_hash=b58decode(converted_content['high']) version=2, content_hash=b58decode(converted_content['high' if content_kind=='video' else 'low'])
).serialize_v2() ).serialize_v2()
unprocessed_encrypted_content.ipfs_cid = ContentId( unprocessed_encrypted_content.ipfs_cid = ContentId(
version=2, content_hash=b58decode(converted_content['low']) version=2, content_hash=b58decode(converted_content['low'])
@@ -199,22 +233,25 @@ async def convert_loop(memory):
**unprocessed_encrypted_content.meta, **unprocessed_encrypted_content.meta,
'converted_content': converted_content 'converted_content': converted_content
} }
session.commit() session.commit()
# Notify user if needed
if not unprocessed_encrypted_content.meta.get('upload_notify_msg_id'): if not unprocessed_encrypted_content.meta.get('upload_notify_msg_id'):
wallet_owner_connection = session.query(WalletConnection).filter( wallet_owner_connection = session.query(WalletConnection).filter(
WalletConnection.wallet_address == unprocessed_encrypted_content.owner_address WalletConnection.wallet_address == unprocessed_encrypted_content.owner_address
).order_by(WalletConnection.id.desc()).first() ).order_by(WalletConnection.id.desc()).first()
if wallet_owner_connection: if wallet_owner_connection:
wallet_owner_user = wallet_owner_connection.user wallet_owner_user = wallet_owner_connection.user
wallet_owner_bot = Wrapped_CBotChat(memory._client_telegram_bot, chat_id=wallet_owner_user.telegram_id, user=wallet_owner_user, db_session=session) bot = Wrapped_CBotChat(
unprocessed_encrypted_content.meta = { memory._client_telegram_bot,
**unprocessed_encrypted_content.meta, chat_id=wallet_owner_user.telegram_id,
'upload_notify_msg_id': await wallet_owner_bot.send_content(session, unprocessed_encrypted_content) 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() session.commit()
async def main_fn(memory): async def main_fn(memory):
make_log("ConvertProcess", "Service started", level="info") make_log("ConvertProcess", "Service started", level="info")
seqno = 0 seqno = 0
+1 -1
View File
@@ -76,7 +76,7 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
).order_by(desc(WalletConnection.id)).first() ).order_by(desc(WalletConnection.id)).first()
wallet_owner_user = wallet_owner_connection.user wallet_owner_user = wallet_owner_connection.user
if wallet_owner_user.telegram_id: if wallet_owner_user.telegram_id:
wallet_owner_bot = Wrapped_CBotChat(memory._client_telegram_bot, chat_id=wallet_owner_user.telegram_id, user=wallet_owner_user, db_session=session) 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( await wallet_owner_bot.send_message(
user.translated('p_licenseWasBought').format( user.translated('p_licenseWasBought').format(
username=user.front_format(), username=user.front_format(),
+7 -5
View File
@@ -3,7 +3,7 @@ from base64 import b64decode
from datetime import datetime, timedelta from datetime import datetime, timedelta
from base58 import b58encode from base58 import b58encode
from sqlalchemy import and_ from sqlalchemy import and_, or_
from tonsdk.boc import Cell from tonsdk.boc import Cell
from tonsdk.utils import Address from tonsdk.utils import Address
@@ -74,12 +74,14 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
# Проверка кошельков пользователей на появление новых NFT, добавление их в базу как неопознанные # Проверка кошельков пользователей на появление новых NFT, добавление их в базу как неопознанные
for user in session.query(User).filter( for user in session.query(User).filter(
User.last_use > datetime.now() - timedelta(minutes=10) User.last_use > datetime.now() - timedelta(hours=4)
).all(): ).order_by(User.updated.asc()).all():
if not user.wallet_address(session): user_wallet_address = user.wallet_address(session)
if not user_wallet_address:
make_log("LicenseIndex", f"User {user.id} has no wallet address", level="info") make_log("LicenseIndex", f"User {user.id} has no wallet address", level="info")
continue continue
make_log("LicenseIndex", f"User {user.id} has wallet address {user_wallet_address}", level="info")
last_updated_licenses = user.meta.get('last_updated_licenses') last_updated_licenses = user.meta.get('last_updated_licenses')
must_skip = last_updated_licenses and (datetime.now() - datetime.fromisoformat(last_updated_licenses)) < timedelta(minutes=1) must_skip = last_updated_licenses and (datetime.now() - datetime.fromisoformat(last_updated_licenses)) < timedelta(minutes=1)
make_log("LicenseIndex", f"User: {user.id}, last_updated_licenses: {last_updated_licenses}, must_skip: {must_skip}", level="info") make_log("LicenseIndex", f"User: {user.id}, last_updated_licenses: {last_updated_licenses}, must_skip: {must_skip}", level="info")
@@ -99,7 +101,7 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
UserContent.type.startswith('nft/'), UserContent.type.startswith('nft/'),
UserContent.updated < (datetime.now() - timedelta(minutes=60)), UserContent.updated < (datetime.now() - timedelta(minutes=60)),
) )
).first() ).order_by(UserContent.updated.asc()).first()
if process_content: if process_content:
make_log("LicenseIndex", f"Syncing content with blockchain: {process_content.id}", level="info") make_log("LicenseIndex", f"Syncing content with blockchain: {process_content.id}", level="info")
try: try:
+7 -4
View File
@@ -127,6 +127,7 @@ async def main_fn(memory):
with db_session() as session: 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): async def process_incoming_transaction(transaction: dict):
transaction_hash = transaction['transaction_id']['hash'] transaction_hash = transaction['transaction_id']['hash']
transaction_lt = str(transaction['transaction_id']['lt']) transaction_lt = str(transaction['transaction_id']['lt'])
@@ -156,7 +157,7 @@ async def main_fn(memory):
in_msg_blockchain_task.status = 'done' in_msg_blockchain_task.status = 'done'
in_msg_blockchain_task.transaction_hash = transaction_hash in_msg_blockchain_task.transaction_hash = transaction_hash
in_msg_blockchain_task.transaction_lt = transaction_lt in_msg_blockchain_task.transaction_lt = transaction_lt
await session.commit() session.commit()
for blockchain_message in [transaction['in_msg']]: for blockchain_message in [transaction['in_msg']]:
try: try:
@@ -174,6 +175,7 @@ async def main_fn(memory):
except BaseException as e: except BaseException as e:
make_log("TON_Daemon", f"Error while getting service wallet transactions: {e}", level="ERROR") 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 ( for blockchain_task in (
session.query(BlockchainTask).filter( session.query(BlockchainTask).filter(
@@ -208,11 +210,12 @@ async def main_fn(memory):
# or sum([int("terminating vm with exit code 36" in e) for e in errors_list]) > 0: # 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") make_log("TON_Daemon", f"Task {blockchain_task.id} done", level="DEBUG")
blockchain_task.status = 'done' blockchain_task.status = 'done'
await session.commit() session.commit()
continue continue
await asyncio.sleep(0.5) await asyncio.sleep(0.5)
await send_status("ton_daemon", f"working: creating new messages (seqno={sw_seqno_value})")
# Создание новых подписей # Создание новых подписей
for blockchain_task in ( for blockchain_task in (
session.query(BlockchainTask).filter(BlockchainTask.status == 'wait').all() session.query(BlockchainTask).filter(BlockchainTask.status == 'wait').all()
@@ -261,9 +264,9 @@ async def main_fn(memory):
blockchain_task.meta = { blockchain_task.meta = {
**blockchain_task.meta, **blockchain_task.meta,
'sign_created': sign_created, 'sign_created': sign_created,
'signed_message': query_boc, 'signed_message': query_boc.hex(),
} }
await session.commit() session.commit()
make_log("TON", f"Created signed message for task {blockchain_task.id}" + '\n' + traceback.format_exc(), level="info") make_log("TON", f"Created signed message for task {blockchain_task.id}" + '\n' + traceback.format_exc(), level="info")
except BaseException as e: except BaseException as e:
make_log("TON", f"Error processing task {blockchain_task.id}: {e}" + '\n' + traceback.format_exc(), level="error") make_log("TON", f"Error processing task {blockchain_task.id}: {e}" + '\n' + traceback.format_exc(), level="error")
+2
View File
@@ -51,6 +51,7 @@ async def create_metadata_for_item(
cover_url: str = None, cover_url: str = None,
authors: list = None, authors: list = None,
hashtags: list = [], hashtags: list = [],
downloadable: bool = False,
) -> StoredContent: ) -> StoredContent:
assert title, "No title provided" assert title, "No title provided"
# assert cover_url, "No cover_url provided" # assert cover_url, "No cover_url provided"
@@ -66,6 +67,7 @@ async def create_metadata_for_item(
# 'value': 'Unknown' # 'value': 'Unknown'
# }, # },
], ],
'downloadable': downloadable,
} }
if cover_url: if cover_url:
item_metadata['image'] = cover_url item_metadata['image'] = cover_url
+2 -28
View File
@@ -131,33 +131,6 @@ class PlayerTemplates:
if cover_content: if cover_content:
# Add thumbnail if cover content is available # Add thumbnail if cover content is available
template_kwargs['thumbnail'] = URLInputFile(cover_content.web_url) template_kwargs['thumbnail'] = URLInputFile(cover_content.web_url)
if self.bot_id == 1:
# Buttons for sharing and opening in app
inline_keyboard_array.append([
{
'text': self.user.translated('shareVideo_button'),
'switch_inline_query': f"Q{user_existing_license.onchain_address}" if user_existing_license else f"C{content.cid.serialize_v2()}",
},
{
'text': self.user.translated('shareLink_button'),
'url': f"https://t.me/share/url?text={urllib.parse.quote(content_share_link['text'])}&url={urllib.parse.quote(content_share_link['url'])}"
}
])
inline_keyboard_array.append([{
'text': self.user.translated('openTrackInApp_button'),
'url': f"https://t.me/{CLIENT_TELEGRAM_BOT_USERNAME}/content?startapp={content.cid.serialize_v2()}"
}])
else:
# Buttons for viewing as a client and opening contract
inline_keyboard_array.append([{
'text': self.user.translated('viewTrackAsClient_button'),
'url': f"https://t.me/{CLIENT_TELEGRAM_BOT_USERNAME}?start=C{content.cid.serialize_v2()}"
}])
inline_keyboard_array.append([{
'text': self.user.translated('openContractPage_button'),
'url': f"https://tonviewer.com/address/{content_meta['item_address']}"
}])
else: else:
local_content = None local_content = None
@@ -168,6 +141,7 @@ class PlayerTemplates:
) )
inline_keyboard_array = [] inline_keyboard_array = []
extra_buttons = [] extra_buttons = []
else: else:
text = content_metadata_json.get('description').strip() text = content_metadata_json.get('description').strip()
@@ -203,7 +177,7 @@ class PlayerTemplates:
await self.delete_message(kmsg.message_id) await self.delete_message(kmsg.message_id)
r = await tg_process_template( r = await tg_process_template(
self, text, message_id=message_id, **template_kwargs, self, text + '\n\n' + f"""<a href="https://t.me/MY_Web3Bot/content?startapp={content.cid.serialize_v2()}"><code>🌐 Открыть на MY</code></a>""", message_id=message_id, **template_kwargs,
keyboard=get_inline_keyboard([*inline_keyboard_array, *extra_buttons]) if inline_keyboard_array else None, keyboard=get_inline_keyboard([*inline_keyboard_array, *extra_buttons]) if inline_keyboard_array else None,
message_type=f'content/{content_type}', message_type=f'content/{content_type}',
message_meta={'content_sha256': content_meta['hash']} if local_content else {}, message_meta={'content_sha256': content_meta['hash']} if local_content else {},
+1 -1
View File
@@ -3,7 +3,7 @@ from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime
from datetime import datetime from datetime import datetime
class PromoAction: class PromoAction(AlchemyBase):
__tablename__ = 'promo_actions' __tablename__ = 'promo_actions'
id = Column(Integer, autoincrement=True, primary_key=True) id = Column(Integer, autoincrement=True, primary_key=True)
+1 -1
View File
@@ -3,7 +3,7 @@ from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime
from datetime import datetime from datetime import datetime
class BlockchainTask: class BlockchainTask(AlchemyBase):
__tablename__ = 'blockchain_tasks' __tablename__ = 'blockchain_tasks'
id = Column(Integer, autoincrement=True, primary_key=True) id = Column(Integer, autoincrement=True, primary_key=True)
+4 -2
View File
@@ -1,3 +1,4 @@
from datetime import datetime
from sqlalchemy import Column, Integer, String, BigInteger, DateTime, JSON from sqlalchemy import Column, Integer, String, BigInteger, DateTime, JSON
from sqlalchemy.orm import relationship from sqlalchemy.orm import relationship
@@ -19,8 +20,9 @@ class User(AlchemyBase, DisplayMixin, TranslationCore, AuthenticationMixin_V1, W
lang_code = Column(String(8), nullable=False, default="en") lang_code = Column(String(8), nullable=False, default="en")
meta = Column(JSON, nullable=False, default={}) meta = Column(JSON, nullable=False, default={})
last_use = Column(DateTime, nullable=False, default=0) last_use = Column(DateTime, nullable=False, default=datetime.utcnow)
created = Column(DateTime, nullable=False, default=0) updated = Column(DateTime, nullable=False, default=datetime.utcnow)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
balances = relationship('UserBalance', back_populates='user') balances = relationship('UserBalance', back_populates='user')
internal_transactions = relationship('InternalTransaction', back_populates='user') internal_transactions = relationship('InternalTransaction', back_populates='user')
+6
View File
@@ -8,6 +8,7 @@ services:
- .env - .env
volumes: volumes:
- /Storage/sqlStorage:/var/lib/mysql - /Storage/sqlStorage:/var/lib/mysql
restart: always
healthcheck: healthcheck:
test: [ "CMD", "healthcheck.sh", "--connect", "--innodb_initialized" ] test: [ "CMD", "healthcheck.sh", "--connect", "--innodb_initialized" ]
interval: 10s interval: 10s
@@ -21,6 +22,7 @@ services:
command: python -m app command: python -m app
env_file: env_file:
- .env - .env
restart: always
links: links:
- maria_db - maria_db
ports: ports:
@@ -36,6 +38,7 @@ services:
build: build:
context: . context: .
dockerfile: Dockerfile dockerfile: Dockerfile
restart: always
command: python -m app indexer command: python -m app indexer
env_file: env_file:
- .env - .env
@@ -53,6 +56,7 @@ services:
context: . context: .
dockerfile: Dockerfile dockerfile: Dockerfile
command: python -m app ton_daemon command: python -m app ton_daemon
restart: always
env_file: env_file:
- .env - .env
links: links:
@@ -69,6 +73,7 @@ services:
context: . context: .
dockerfile: Dockerfile dockerfile: Dockerfile
command: python -m app license_index command: python -m app license_index
restart: always
env_file: env_file:
- .env - .env
links: links:
@@ -85,6 +90,7 @@ services:
context: . context: .
dockerfile: Dockerfile dockerfile: Dockerfile
command: python -m app convert_process command: python -m app convert_process
restart: always
env_file: env_file:
- .env - .env
links: links:
Binary file not shown.
+5 -2
View File
@@ -148,8 +148,11 @@ msgid "buyTrackListenLicense_button"
msgstr "💶 Купить лицензию ({price} TON)" msgstr "💶 Купить лицензию ({price} TON)"
#: app/core/models/_telegram/templates/player.py:117 #: app/core/models/_telegram/templates/player.py:117
msgid "p_playerContext_preview" msgid "p_playerAudioContext_preview"
msgstr "<i>Это демонстрационная версия трека. Чтобы прослушать полную версию, необходимо приобрести лицензию 🎧.</i>" msgstr "<i>Превью аудиозаписи. Полная версия доступна только по лицензии 🎧</i>"
msgid "p_playerVideoContext_preview"
msgstr "<i>Превью видео. Полная версия доступна только по лицензии 🎬</i>"
#: app/client_bot/routers/content.py:32 #: app/client_bot/routers/content.py:32
msgid "error_contentNotFound" msgid "error_contentNotFound"
+2
View File
@@ -15,4 +15,6 @@ aiofiles==23.2.1
pydub==0.25.1 pydub==0.25.1
pillow==10.2.0 pillow==10.2.0
ffmpeg-python==0.2.0 ffmpeg-python==0.2.0
python-magic==0.4.27