55 Commits
Author SHA1 Message Date
user 846e32c5b1 edit platform deployment 2025-08-25 14:36:10 +03:00
user e9e2f25f4d fix lazy loading 2025-08-25 14:35:20 +03:00
user 608881b5d8 fix misprint 2025-08-25 14:01:34 +03:00
user 3747329b1e fix misprint 2025-08-25 13:50:05 +03:00
user 4d5318b5d4 new deploy logic 2025-08-25 13:34:12 +03:00
user 3e6d0b93cb update versions contracts 2025-08-25 13:22:30 +03:00
user 2bd6e30b38 edit platfrom contract deployment 2025-08-25 12:25:08 +03:00
user 45374987e7 secrets fix 2025-08-25 11:52:29 +03:00
user 7d920907cc fix secrets read stuck 2025-08-24 19:39:28 +03:00
user 4401916104 fix db connections 2025-08-24 17:17:58 +03:00
user 3a6f787a78 update platform code contract 2025-08-24 16:48:23 +03:00
user d67135849c fix converter module path 2025-08-24 14:13:14 +03:00
user 4da4cd1526 fix db init 2025-08-24 13:23:05 +03:00
user f562dc8ed7 fix lazy_loading 2025-08-24 13:11:55 +03:00
user 5d41d33c6e fix syntax err 2025-08-23 21:42:14 +03:00
user 61e85baf08 extend logs 2025-08-23 21:12:37 +03:00
user 695969f015 better logs 2025-08-23 18:54:53 +03:00
user 82758fb11a postgres default fix 2025-08-23 13:15:10 +03:00
user b28e561a5f edit platform contract 2025-08-23 12:24:14 +03:00
user cf64ddaaa5 fix misprint 2025-08-22 19:36:20 +03:00
user 79165b49b5 startup errors fix #2 2025-08-22 19:09:38 +03:00
user 4cca40a626 startup errors fix #1 2025-08-22 14:12:45 +03:00
user e51bb86dc0 mariadb -> postgres 2025-08-22 14:04:21 +03:00
user 21964fa986 fix player ui 2025-06-01 12:17:35 +03:00
user 590afd2475 text pre-filtering 2025-06-01 11:58:51 +03:00
user a266c8b710 fix misprint 2025-06-01 08:13:12 +03:00
user b5a9437c05 fix player ui 2025-06-01 00:04:26 +03:00
user 58eca166db edit player ui globally 2025-05-31 23:45:02 +03:00
user 27dc827880 new player version 2025-05-31 01:55:16 +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
user d852913692 free upload & player 2025-03-13 18:34:08 +03:00
user 862683b36e add topup highload before deploy 2025-03-13 16:33:37 +03:00
user 446bd74464 fix ton daemon 2025-03-13 16:23:41 +03:00
user b8a5cf4965 add highload wallet 2025-03-13 15:46:48 +03:00
60 changed files with 1729 additions and 845 deletions

No files matched your search

+3 -2
View File
@@ -12,7 +12,8 @@ RUN apt-get update && apt-get install -y \
ca-certificates \
curl \
gnupg \
lsb-release && \
lsb-release \
ffmpeg && \
install -m 0755 -d /etc/apt/keyrings && \
curl -fsSL https://download.docker.com/linux/debian/gpg -o /etc/apt/keyrings/docker.asc && \
chmod a+r /etc/apt/keyrings/docker.asc && \
@@ -21,7 +22,7 @@ RUN apt-get update && apt-get install -y \
apt-get update && \
apt-get install -y docker-ce-cli
RUN apt-get install -y ffmpeg
RUN apt-get install libmagic1 -y
CMD ["python", "app"]
+36 -9
View File
@@ -12,17 +12,12 @@ try:
except BaseException:
pass
from app.core._utils.create_maria_tables import create_maria_tables
from app.core._utils.create_maria_tables import create_db_tables
from app.core.storage import engine
if startup_target == '__main__':
create_maria_tables(engine)
else:
if startup_target != '__main__':
# Background services get a short delay before startup
time.sleep(7)
from app.api import app
from app.bot import dp as uploader_bot_dp
from app.client_bot import dp as client_bot_dp
from app.core._config import SANIC_PORT, MYSQL_URI, PROJECT_HOST
from app.core.logger import make_log
if int(os.getenv("SANIC_MAINTENANCE", '0')) == 1:
@@ -52,7 +47,11 @@ async def execute_queue(app):
make_log(None, f"Application normally started. HTTP port: {SANIC_PORT}")
make_log(None, f"Telegram bot: https://t.me/{telegram_bot_username}")
make_log(None, f"Client Telegram bot: https://t.me/{client_telegram_bot_username}")
make_log(None, f"MariaDB host: {MYSQL_URI.split('@')[1].split('/')[0].replace('/', '')}")
try:
_db_host = DATABASE_URL.split('@')[1].split('/')[0].replace('/', '')
except Exception:
_db_host = 'postgres://'
make_log(None, f"PostgreSQL host: {_db_host}")
make_log(None, f"API host: {PROJECT_HOST}")
while True:
try:
@@ -81,12 +80,39 @@ async def execute_queue(app):
if __name__ == '__main__':
main_memory = Memory()
if startup_target == '__main__':
# Defer heavy imports to avoid side effects in background services
# Mark this process as the primary node for seeding/config init
os.environ.setdefault('NODE_ROLE', 'primary')
# Create DB tables synchronously before importing HTTP app to satisfy _secrets
try:
from sqlalchemy import create_engine
from app.core.models import AlchemyBase # imports all models
db_url = os.environ.get('DATABASE_URL')
if not db_url:
raise RuntimeError('DATABASE_URL is not set')
# Normalize to sync driver
if '+asyncpg' in db_url:
db_url_sync = db_url.replace('+asyncpg', '+psycopg2')
else:
db_url_sync = db_url
sync_engine = create_engine(db_url_sync, pool_pre_ping=True)
AlchemyBase.metadata.create_all(sync_engine)
except Exception as e:
make_log('Startup', f'DB sync init failed: {e}', level='error')
from app.api import app
from app.bot import dp as uploader_bot_dp
from app.client_bot import dp as client_bot_dp
from app.core._config import SANIC_PORT, PROJECT_HOST, DATABASE_URL
app.ctx.memory = main_memory
for _target in [uploader_bot_dp, client_bot_dp]:
_target._s_memory = app.ctx.memory
app.ctx.memory._app = app
# Ensure DB schema exists using the same event loop as Sanic (idempotent)
app.add_task(create_db_tables(engine))
app.add_task(execute_queue(app))
app.add_task(queue_daemon(app))
app.add_task(uploader_bot_dp.start_polling(app.ctx.memory._telegram_bot))
@@ -126,6 +152,7 @@ if __name__ == '__main__':
loop = asyncio.get_event_loop()
try:
# Background services no longer perform schema initialization
loop.run_until_complete(wrapped_startup_fn(main_memory))
except BaseException as e:
make_log(startup_target[0].upper() + startup_target[1:], f"Error: {e}" + '\n' + str(traceback.format_exc()),
+55 -7
View File
@@ -1,6 +1,8 @@
import traceback
from sanic import Sanic, response
from uuid import uuid4
import traceback as _traceback
from app.core.logger import make_log
@@ -12,8 +14,8 @@ app.register_middleware(attach_user_to_request, "request")
app.register_middleware(close_db_session, "response")
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
from app.api.routes.auth import s_api_v1_auth_twa
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, 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.node_storage import s_api_v1_storage_post, s_api_v1_storage_get, \
s_api_v1_storage_decode_cid
@@ -29,6 +31,7 @@ app.add_route(s_index, "/", methods=["GET", "OPTIONS"])
app.add_route(s_favicon, "/favicon.ico", methods=["GET", "OPTIONS"])
app.add_route(s_api_v1_node, "/api/v1/node", methods=["GET", "OPTIONS"])
app.add_route(s_api_v1_node_friendly, "/api/v1/nodeFriendly", methods=["GET", "OPTIONS"])
app.add_route(s_api_system_version, "/api/system.version", methods=["GET", "OPTIONS"])
app.add_route(s_api_system_send_status, "/api/system.sendStatus", methods=["POST", "OPTIONS"])
@@ -36,6 +39,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_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_logout, "/api/v1/tonconnect.logout", methods=["POST", "OPTIONS"])
@@ -60,18 +65,61 @@ app.add_route(s_api_v1_5_content_list, "/api/v1.5/content.list", methods=["GET",
@app.exception(BaseException)
async def s_handle_exception(request, exception):
response_buffer = response.json({"error": "An internal server error occurred"}, status=500)
# Correlate error to request
session_id = getattr(request.ctx, 'session_id', None) or uuid4().hex[:16]
error_id = uuid4().hex[:8]
status = 500
code = type(exception).__name__
message = "Internal HTTP Error"
try:
raise exception
except AssertionError as e:
response_buffer = response.json({"error": str(e)}, status=400)
status = 400
code = 'AssertionError'
message = str(e) or 'Bad Request'
except BaseException as e:
make_log("sanic_exception", f"Exception: {e}" + '\n' + str(traceback.format_exc()), level='error')
# keep default 500, but expose exception message to aid debugging
message = str(e) or message
# Build structured log with full context and traceback
try:
tb = _traceback.format_exc()
user_id = getattr(getattr(request.ctx, 'user', None), 'id', None)
log_ctx = {
'sid': session_id,
'eid': error_id,
'path': request.path,
'method': request.method,
'query': dict(request.args) if hasattr(request, 'args') else {},
'user_id': user_id,
'remote': (request.headers.get('X-Forwarded-For') or request.remote_addr or request.ip),
'code': code,
'message': message,
'traceback': tb,
}
make_log('http_exception', 'API exception', level='error', **log_ctx)
except BaseException:
pass
# Return enriched error response for the client
payload = {
'error': True,
'code': code,
'message': message,
'session_id': session_id,
'error_id': error_id,
'path': request.path,
'method': request.method,
}
response_buffer = response.json(payload, status=status)
response_buffer = await close_db_session(request, response_buffer)
response_buffer.headers["Access-Control-Allow-Origin"] = "*"
response_buffer.headers["Access-Control-Allow-Methods"] = "GET, POST, OPTIONS"
response_buffer.headers["Access-Control-Allow-Headers"] = "Origin, Content-Type, Accept, Authorization, Referer, User-Agent, Sec-Fetch-Dest, Sec-Fetch-Mode, Sec-Fetch-Site"
response_buffer.headers["Access-Control-Allow-Headers"] = "Origin, Content-Type, Accept, Authorization, Referer, User-Agent, Sec-Fetch-Dest, Sec-Fetch-Mode, Sec-Fetch-Site, x-request-id"
response_buffer.headers["Access-Control-Allow-Credentials"] = "true"
response_buffer.headers["X-Session-Id"] = session_id
response_buffer.headers["X-Error-Id"] = error_id
return response_buffer
+81 -15
View File
@@ -1,5 +1,6 @@
from base58 import b58decode
from sanic import response as sanic_response
from uuid import uuid4
from app.core._crypto.signer import Signer
from app.core._secrets import hot_seed
@@ -8,15 +9,25 @@ from app.core.models.keys import KnownKey
from app.core.models._telegram.wrapped_bot import Wrapped_CBotChat
from app.core.models.user_activity import UserActivity
from app.core.models.user import User
from app.core.storage import Session
from sqlalchemy import select
from app.core.storage import new_session
from datetime import datetime, timedelta
from app.core.log_context import (
ctx_session_id, ctx_user_id, ctx_method, ctx_path, ctx_remote
)
def attach_headers(response):
def attach_headers(response, request=None):
response.headers["Access-Control-Allow-Origin"] = "*"
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-Credentials"] = "true"
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, x-request-id"
# response.headers["Access-Control-Allow-Credentials"] = "true"
try:
sid = getattr(request.ctx, 'session_id', None) if request else None
if sid:
response.headers["X-Session-Id"] = sid
except BaseException:
pass
return response
@@ -30,7 +41,8 @@ async def try_authorization(request):
make_log("auth", "Invalid token length", level="warning")
return
known_key = request.ctx.db_session.query(KnownKey).filter(KnownKey.seed == token).first()
result = await request.ctx.db_session.execute(select(KnownKey).where(KnownKey.seed == token))
known_key = result.scalars().first()
if not known_key:
make_log("auth", "Unknown key", level="warning")
return
@@ -58,7 +70,8 @@ async def try_authorization(request):
make_log("auth", f"User ID mismatch: {known_key.meta.get('I_user_id', -1)} != {user_id}", level="warning")
return
user = request.ctx.db_session.query(User).filter(User.id == known_key.meta['I_user_id']).first()
result = await request.ctx.db_session.execute(select(User).where(User.id == known_key.meta['I_user_id']))
user = result.scalars().first()
if not user:
make_log("auth", "No user from key", level="warning")
return
@@ -118,7 +131,14 @@ async def save_activity(request):
pass
try:
activity_meta["headers"] = dict(request.headers)
# Sanitize sensitive headers
headers = dict(request.headers)
for hk in list(headers.keys()):
if str(hk).lower() in [
'authorization', 'cookie', 'x-service-signature', 'x-message-hash'
]:
headers[hk] = '<redacted>'
activity_meta["headers"] = headers
except:
pass
@@ -127,23 +147,51 @@ async def save_activity(request):
meta=activity_meta,
user_id=request.ctx.user.id if request.ctx.user else None,
user_ip=activity_meta.get("ip", "0.0.0.0"),
created=datetime.now()
created=datetime.utcnow()
)
request.ctx.db_session.add(new_user_activity)
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
async def attach_user_to_request(request):
if request.method == 'OPTIONS':
return attach_headers(sanic_response.text("OK"))
return attach_headers(sanic_response.text("OK"), request)
request.ctx.db_session = Session()
request.ctx.db_session = new_session()
request.ctx.verified_hash = None
request.ctx.user = None
request.ctx.user_key = None
request.ctx.user_uploader_wrapper = Wrapped_CBotChat(request.app.ctx.memory._telegram_bot, db_session=request.ctx.db_session)
request.ctx.user_client_wrapper = Wrapped_CBotChat(request.app.ctx.memory._client_telegram_bot, db_session=request.ctx.db_session)
# Correlation/session id for this request: prefer proxy-provided X-Request-ID
incoming_req_id = request.headers.get('X-Request-Id') or request.headers.get('X-Request-ID')
request.ctx.session_id = (incoming_req_id or uuid4().hex)[:32]
# Populate contextvars for automatic logging context
try:
ctx_session_id.set(request.ctx.session_id)
ctx_method.set(request.method)
ctx_path.set(request.path)
_remote = (request.headers.get('X-Forwarded-For') or request.remote_addr or request.ip)
if _remote and isinstance(_remote, str) and ',' in _remote:
_remote = _remote.split(',')[0].strip()
ctx_remote.set(_remote)
except BaseException:
pass
try:
make_log(
"HTTP",
f"Request start sid={request.ctx.session_id} {request.method} {request.path}",
level='info'
)
except BaseException:
pass
await try_authorization(request)
# Update user_id in context after auth
try:
if request.ctx.user and request.ctx.user.id:
ctx_user_id.set(request.ctx.user.id)
except BaseException:
pass
await save_activity(request)
await try_service_authorization(request)
@@ -153,16 +201,34 @@ async def close_request_handler(request, response):
response = sanic_response.text("OK")
try:
request.ctx.db_session.close()
except BaseException as e:
await request.ctx.db_session.close()
except BaseException:
pass
response = attach_headers(response)
try:
make_log(
"HTTP",
f"Request end sid={getattr(request.ctx, 'session_id', None)} {request.method} {request.path} status={getattr(response, 'status', None)}",
level='info'
)
except BaseException:
pass
response = attach_headers(response, request)
return request, response
async def close_db_session(request, response):
request, response = await close_request_handler(request, response)
response = attach_headers(response)
response = attach_headers(response, request)
# Clear contextvars
try:
ctx_session_id.set(None)
ctx_user_id.set(None)
ctx_method.set(None)
ctx_path.set(None)
ctx_remote.set(None)
except BaseException:
pass
return response
+136 -16
View File
@@ -3,6 +3,7 @@ from datetime import datetime
import traceback
from sanic import response
from sqlalchemy import and_, select, func
from tonsdk.boc import begin_cell, begin_dict
from tonsdk.utils import Address
@@ -18,6 +19,8 @@ from app.core.models.content.user_content import UserContent
from app.core.models.node_storage import StoredContent
from app.core.models._telegram import Wrapped_CBotChat
from app.core._keyboards import get_inline_keyboard
from app.core.models.promo import PromoAction
from app.core.models.tasks import BlockchainTask
def valid_royalty_params(royalty_params):
@@ -58,9 +61,9 @@ async def s_api_v1_blockchain_send_new_content_message(request):
assert not err, f"Invalid content CID"
# Поиск исходного файла загруженного
decrypted_content = request.ctx.db_session.query(StoredContent).filter(
StoredContent.hash == decrypted_content_cid.content_hash_b58
).first()
decrypted_content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == decrypted_content_cid.content_hash_b58)
)).scalars().first()
assert decrypted_content, "No content locally found"
assert decrypted_content.type == "local/content_bin", "Invalid content type"
@@ -71,23 +74,24 @@ async def s_api_v1_blockchain_send_new_content_message(request):
if request.json['image']:
image_content_cid, err = resolve_content(request.json['image'])
assert not err, f"Invalid image CID"
image_content = request.ctx.db_session.query(StoredContent).filter(
StoredContent.hash == image_content_cid.content_hash_b58
).first()
image_content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == image_content_cid.content_hash_b58)
)).scalars().first()
assert image_content, "No image locally found"
else:
image_content_cid = None
image_content = None
content_title = f"{', '.join(request.json['authors'])} - {request.json['title']}" if request.json['authors'] else request.json['title']
content_title = f"{', '.join(request.json['authors'])} – {request.json['title']}" if request.json['authors'] else request.json['title']
metadata_content = await create_metadata_for_item(
request.ctx.db_session,
title=content_title,
cover_url=f"{PROJECT_HOST}/api/v1.5/storage/{image_content_cid.serialize_v2()}" if image_content_cid else None,
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)
@@ -101,6 +105,99 @@ async def s_api_v1_blockchain_send_new_content_message(request):
)
i += 1
_cnt = (await request.ctx.db_session.execute(
select(func.count()).select_from(PromoAction).where(
and_(
PromoAction.user_internal_id == request.ctx.user.id,
PromoAction.action_type == 'freeUpload'
)
)
)).scalar()
promo_free_upload_available = 3 - int(_cnt or 0)
has_pending_task = (await request.ctx.db_session.execute(
select(BlockchainTask).where(
and_(BlockchainTask.user_id == request.ctx.user.id, BlockchainTask.status != 'done')
)
)).scalars().first()
if has_pending_task:
make_log("Blockchain", f"User {request.ctx.user.id} already has a pending task", level='warning')
promo_free_upload_available = 0
make_log("Blockchain", f"User {request.ctx.user.id} has {promo_free_upload_available} free uploads available", level='info')
if promo_free_upload_available > 0:
promo_action = PromoAction(
user_id = str(request.ctx.user.id),
user_internal_id=request.ctx.user.id,
action_type='freeUpload',
action_ref=str(encrypted_content_cid.content_hash),
created=datetime.now()
)
request.ctx.db_session.add(promo_action)
blockchain_task = BlockchainTask(
destination=platform.address.to_string(1, 1, 1),
amount=str(int(0.03 * 10 ** 9)),
payload=b64encode(
begin_cell()
.store_uint(0x5491d08c, 32)
.store_uint(int.from_bytes(encrypted_content_cid.content_hash, "big", signed=False), 256)
.store_address(Address(await request.ctx.user.wallet_address_async(request.ctx.db_session)))
.store_ref(
begin_cell()
.store_ref(
begin_cell()
.store_coins(int(0))
.store_coins(int(0))
.store_coins(int(request.json['price']))
.end_cell()
)
.store_maybe_ref(royalties_dict.end_dict())
.store_uint(0, 1)
.end_cell()
)
.store_ref(
begin_cell()
.store_ref(
begin_cell()
.store_bytes(f"{PROJECT_HOST}/api/v1.5/storage/{metadata_content.cid.serialize_v2(include_accept_type=True)}".encode())
.end_cell()
)
.store_ref(
begin_cell()
.store_ref(begin_cell().store_bytes(f"{encrypted_content_cid.serialize_v2()}".encode()).end_cell())
.store_ref(begin_cell().store_bytes(f"{image_content_cid.serialize_v2() if image_content_cid else ''}".encode()).end_cell())
.store_ref(begin_cell().store_bytes(f"{metadata_content.cid.serialize_v2()}".encode()).end_cell())
.end_cell()
)
.end_cell()
)
.end_cell().to_boc(False)
).decode(),
epoch=None, seqno=None,
created = datetime.now(),
status='wait',
user_id = request.ctx.user.id
)
request.ctx.db_session.add(blockchain_task)
await request.ctx.db_session.commit()
await request.ctx.user_uploader_wrapper.send_message(
request.ctx.user.translated('p_uploadContentTxPromo').format(
title=content_title,
free_count=(promo_free_upload_available - 1)
), message_type='hint', message_meta={
'encrypted_content_hash': b58encode(encrypted_content_cid.content_hash).decode(),
'hint_type': 'uploadContentTxRequested'
}
)
return response.json({
'address': "free",
'amount': str(int(0.03 * 10 ** 9)),
'payload': ""
})
await request.ctx.user_uploader_wrapper.send_message(
request.ctx.user.translated('p_uploadContentTxRequested').format(
title=content_title,
@@ -117,6 +214,7 @@ async def s_api_v1_blockchain_send_new_content_message(request):
begin_cell()
.store_uint(0x5491d08c, 32)
.store_uint(int.from_bytes(encrypted_content_cid.content_hash, "big", signed=False), 256)
.store_uint(0, 2)
.store_ref(
begin_cell()
.store_ref(
@@ -164,15 +262,37 @@ async def s_api_v1_blockchain_send_purchase_content_message(request):
assert field_key in request.json, f"No {field_key} provided"
assert field_value(request.json[field_key]), f"Invalid {field_key} provided"
license_exist = request.ctx.db_session.query(UserContent).filter_by(
onchain_address=request.json['content_address'],
).first()
if license_exist:
r_content = StoredContent.from_cid(request.ctx.db_session, license_exist.content.cid.serialize_v2())
else:
r_content = StoredContent.from_cid(request.ctx.db_session, request.json['content_address'])
if not (await request.ctx.user.wallet_address_async(request.ctx.db_session)):
return response.json({"error": "No wallet address provided"}, status=400)
content = r_content.open_content(request.ctx.db_session)
from sqlalchemy import select
license_exist = (await request.ctx.db_session.execute(select(UserContent).where(
UserContent.onchain_address == request.json['content_address']
))).scalars().first()
if license_exist:
from app.core.content.content_id import ContentId
_cid = ContentId.deserialize(license_exist.content.cid.serialize_v2())
r_content = (await request.ctx.db_session.execute(select(StoredContent).where(StoredContent.hash == _cid.content_hash_b58))).scalars().first()
else:
from app.core.content.content_id import ContentId
_cid = ContentId.deserialize(request.json['content_address'])
r_content = (await request.ctx.db_session.execute(select(StoredContent).where(StoredContent.hash == _cid.content_hash_b58))).scalars().first()
async def open_content_async(session, sc: StoredContent):
if not sc.encrypted:
decrypted = sc
encrypted = (await session.execute(select(StoredContent).where(StoredContent.decrypted_content_id == sc.id))).scalars().first()
else:
encrypted = sc
decrypted = (await session.execute(select(StoredContent).where(StoredContent.id == sc.decrypted_content_id))).scalars().first()
assert decrypted and encrypted, "Can't open content"
ctype = decrypted.json_format().get('content_type', 'application/x-binary')
try:
content_type = ctype.split('/')[0]
except Exception:
content_type = 'application'
return {'encrypted_content': encrypted, 'decrypted_content': decrypted, 'content_type': content_type}
content = await open_content_async(request.ctx.db_session, r_content)
licenses_cost = content['encrypted_content'].json_format()['license']
assert request.json['license_type'] in licenses_cost
+26 -4
View File
@@ -6,6 +6,7 @@ from base58 import b58encode, b58decode
from sanic import response
from app.core.models.node_storage import StoredContent
from sqlalchemy import select
from app.core._blockchain.ton.platform import platform
from app.core._crypto.signer import Signer
from app.core._secrets import hot_pubkey, service_wallet, hot_seed
@@ -19,10 +20,10 @@ def get_git_info():
async def s_api_v1_node(request): # /api/v1/node
last_known_index = request.ctx.db_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_obj = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.onchain_index != None).order_by(StoredContent.onchain_index.desc())
)).scalars().first()
last_known_index = last_known_index_obj.onchain_index if last_known_index_obj else 0
last_known_index = max(last_known_index, 0)
return response.json({
'id': b58encode(hot_pubkey).decode(),
@@ -38,6 +39,27 @@ async def s_api_v1_node(request): # /api/v1/node
}
})
async def s_api_v1_node_friendly(request):
last_known_index_obj = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.onchain_index != None).order_by(StoredContent.onchain_index.desc())
)).scalars().first()
last_known_index = last_known_index_obj.onchain_index if last_known_index_obj else 0
last_known_index = max(last_known_index, 0)
response_plain_text = f"""
Node address: {service_wallet.address.to_string(1, 1, 1)}
Node ID: {b58encode(hot_pubkey).decode()}
Master address: {platform.address.to_string(1, 1, 1)}
Indexer height: {last_known_index}
Services:
"""
for service_key, service in request.app.ctx.memory.known_states.items():
response_plain_text += f"""
{service_key}:
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}
"""
return response.text(response_plain_text, content_type='text/plain')
async def s_api_system_send_status(request):
if not request.json:
+86 -11
View File
@@ -1,4 +1,5 @@
from datetime import datetime
from uuid import uuid4
from aiogram.utils.web_app import safe_parse_webapp_init_data
from sanic import response
@@ -36,7 +37,9 @@ async def s_api_v1_auth_twa(request):
make_log("auth", "Invalid TWA data", level="warning")
return response.json({"error": "Invalid TWA data"}, status=401)
known_user = request.ctx.db_session.query(User).filter(User.telegram_id == twa_data.user.id).first()
known_user = (await request.ctx.db_session.execute(
select(User).where(User.telegram_id == twa_data.user.id)
)).scalars().first()
if not known_user:
new_user = User(
telegram_id=twa_data.user.id,
@@ -51,9 +54,11 @@ async def s_api_v1_auth_twa(request):
created=datetime.now()
)
request.ctx.db_session.add(new_user)
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
known_user = request.ctx.db_session.query(User).filter(User.telegram_id == twa_data.user.id).first()
known_user = (await request.ctx.db_session.execute(
select(User).where(User.telegram_id == twa_data.user.id)
)).scalars().first()
assert known_user, "User not created"
new_user_key = await known_user.create_api_token_v1(request.ctx.db_session, "USER_API_V1")
@@ -64,12 +69,12 @@ async def s_api_v1_auth_twa(request):
wallet_info.account = Account.from_dict(auth_data['ton_proof']['account'])
wallet_info.ton_proof = TonProof.from_dict({'proof': auth_data['ton_proof']['ton_proof']})
connection_payload = auth_data['ton_proof']['ton_proof']['payload']
known_payload = (request.ctx.db_session.execute(select(KnownKey).where(KnownKey.seed == connection_payload))).scalars().first()
known_payload = (await request.ctx.db_session.execute(select(KnownKey).where(KnownKey.seed == connection_payload))).scalars().first()
assert known_payload, "Unknown payload"
assert known_payload.meta['I_user_id'] == known_user.id, "Invalid user_id"
assert wallet_info.check_proof(connection_payload), "Invalid proof"
for known_connection in (request.ctx.db_session.execute(select(WalletConnection).where(
for known_connection in (await request.ctx.db_session.execute(select(WalletConnection).where(
and_(
WalletConnection.user_id == known_user.id,
WalletConnection.network == 'ton'
@@ -77,7 +82,7 @@ async def s_api_v1_auth_twa(request):
))).scalars().all():
known_connection.invalidated = True
for other_connection in (request.ctx.db_session.execute(select(WalletConnection).where(
for other_connection in (await request.ctx.db_session.execute(select(WalletConnection).where(
WalletConnection.wallet_address == Address(wallet_info.account.address).to_string(1, 1, 1)
))).scalars().all():
other_connection.invalidated = True
@@ -98,24 +103,94 @@ async def s_api_v1_auth_twa(request):
without_pk=False
)
request.ctx.db_session.add(new_connection)
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
except BaseException as e:
make_log("auth", f"Invalid ton_proof: {e}", level="warning")
return response.json({"error": "Invalid ton_proof"}, status=400)
ton_connection = (request.ctx.db_session.execute(select(WalletConnection).where(
ton_connection = (await request.ctx.db_session.execute(select(WalletConnection).where(
and_(
WalletConnection.user_id == known_user.id,
WalletConnection.network == 'ton',
WalletConnection.invalidated == False
)
))).scalars().first()
).order_by(WalletConnection.created.desc()))).scalars().first()
known_user.last_use = datetime.now()
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
return response.json({
'user': known_user.json_format(),
'connected_wallet': ton_connection.json_format() if ton_connection else None,
'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 = (await 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 = (await db_session.execute(select(WalletConnection).where(
and_(
WalletConnection.user_id == user.id,
WalletConnection.wallet_address == canonical_address
)
))).scalars().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)
await db_session.commit()
return response.empty(status=200)
+79 -35
View File
@@ -1,5 +1,6 @@
from datetime import datetime, timedelta
from sanic import response
from sqlalchemy import select, and_, func
from aiogram import Bot, types
from sqlalchemy import and_
from app.core.logger import make_log
@@ -22,13 +23,20 @@ async def s_api_v1_content_list(request):
store = request.args.get('store', 'local')
assert store in ('local', 'onchain'), "Invalid store"
content_list = request.ctx.db_session.query(StoredContent).filter(
stmt = (
select(StoredContent)
.where(
StoredContent.type.like(store + '%'),
StoredContent.disabled == False
).order_by(StoredContent.created.desc()).offset(offset).limit(limit)
make_log("Content", f"Listed {content_list.count()} contents", level='info')
)
.order_by(StoredContent.created.desc())
.offset(offset)
.limit(limit)
)
rows = (await request.ctx.db_session.execute(stmt)).scalars().all()
make_log("Content", f"Listed {len(rows)} contents", level='info')
result = {}
for content in content_list.all():
for content in rows:
content_json = content.json_format()
result[content_json["cid"]] = content_json
@@ -38,23 +46,41 @@ async def s_api_v1_content_list(request):
async def s_api_v1_content_view(request, content_address: str):
# content_address can be CID or TON address
license_exist = request.ctx.db_session.query(UserContent).filter_by(
onchain_address=content_address,
).first()
license_exist = (await request.ctx.db_session.execute(
select(UserContent).where(UserContent.onchain_address == content_address)
)).scalars().first()
if license_exist:
content_address = license_exist.content.cid.serialize_v2()
r_content = StoredContent.from_cid(request.ctx.db_session, content_address)
content = r_content.open_content(request.ctx.db_session)
from app.core.content.content_id import ContentId
cid = ContentId.deserialize(content_address)
r_content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == cid.content_hash_b58)
)).scalars().first()
async def open_content_async(session, sc: StoredContent):
if not sc.encrypted:
decrypted = sc
encrypted = (await session.execute(select(StoredContent).where(StoredContent.decrypted_content_id == sc.id))).scalars().first()
else:
encrypted = sc
decrypted = (await session.execute(select(StoredContent).where(StoredContent.id == sc.decrypted_content_id))).scalars().first()
assert decrypted and encrypted, "Can't open content"
ctype = decrypted.json_format().get('content_type', 'application/x-binary')
try:
content_type = ctype.split('/')[0]
except Exception:
content_type = 'application'
return {'encrypted_content': encrypted, 'decrypted_content': decrypted, 'content_type': content_type}
content = await open_content_async(request.ctx.db_session, r_content)
opts = {
'content_type': content['content_type'], # возможно с ошибками, нужно переделать на ffprobe
'content_address': content['encrypted_content'].meta.get('item_address', '')
}
if content['encrypted_content'].key_id:
known_key = request.ctx.db_session.query(KnownKey).filter(
KnownKey.id == content['encrypted_content'].key_id
).first()
known_key = (await request.ctx.db_session.execute(
select(KnownKey).where(KnownKey.id == content['encrypted_content'].key_id)
)).scalars().first()
if known_key:
opts['key_hash'] = known_key.seed_hash # нахер не нужно на данный момент
@@ -64,22 +90,23 @@ async def s_api_v1_content_view(request, content_address: str):
have_access = False
if request.ctx.user:
user_wallet_address = request.ctx.user.wallet_address(request.ctx.db_session)
user_wallet_address = await request.ctx.user.wallet_address_async(request.ctx.db_session)
have_access = (
(content['encrypted_content'].owner_address == user_wallet_address)
or bool(request.ctx.db_session.query(UserContent).filter_by(owner_address=user_wallet_address, status='active',
content_id=content['encrypted_content'].id).first()) \
or bool(request.ctx.db_session.query(StarsInvoice).filter(
or bool((await request.ctx.db_session.execute(select(UserContent).where(
and_(UserContent.owner_address == user_wallet_address, UserContent.status == 'active', UserContent.content_id == content['encrypted_content'].id)
))).scalars().first()) \
or bool((await request.ctx.db_session.execute(select(StarsInvoice).where(
and_(
StarsInvoice.user_id == request.ctx.user.id,
StarsInvoice.content_hash == content['encrypted_content'].hash,
StarsInvoice.paid == True
)
).first())
))).scalars().first())
)
if not have_access:
current_star_rate = ServiceConfig(request.ctx.db_session).get('live_tonPerStar', [0, 0])[0]
current_star_rate = (await ServiceConfig(request.ctx.db_session).get('live_tonPerStar', [0, 0]))[0]
if current_star_rate < 0:
current_star_rate = 0.00000001
@@ -88,14 +115,14 @@ async def s_api_v1_content_view(request, content_address: str):
stars_cost = 2
invoice_id = f"access_{uuid.uuid4().hex}"
exist_invoice = request.ctx.db_session.query(StarsInvoice).filter(
exist_invoice = (await request.ctx.db_session.execute(select(StarsInvoice).where(
and_(
StarsInvoice.user_id == request.ctx.user.id,
StarsInvoice.created > datetime.now() - timedelta(minutes=25),
StarsInvoice.amount == stars_cost,
StarsInvoice.content_hash == content['encrypted_content'].hash,
)
).first()
))).scalars().first()
if exist_invoice:
invoice_url = exist_invoice.invoice_url
else:
@@ -119,7 +146,7 @@ async def s_api_v1_content_view(request, content_address: str):
invoice_url=invoice_url
)
)
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
except BaseException as e:
make_log("Content", f"Can't create invoice link: {e}", level='warning')
@@ -142,23 +169,33 @@ async def s_api_v1_content_view(request, content_address: str):
if have_access:
user_content_option = 'low' # TODO: подключать high если человек внезапно меломан
converted_content = request.ctx.db_session.query(StoredContent).filter(
converted_content = (await request.ctx.db_session.execute(select(StoredContent).where(
StoredContent.hash == converted_content[user_content_option]
).first()
))).scalars().first()
if converted_content:
display_options['content_url'] = converted_content.web_url
opts['content_ext'] = converted_content.filename.split('.')[-1]
content_meta = content['encrypted_content'].json_format()
content_metadata = StoredContent.from_cid(request.ctx.db_session, content_meta.get('metadata_cid') or None)
from app.core.content.content_id import ContentId
_mcid = content_meta.get('metadata_cid') or None
content_metadata = None
if _mcid:
_cid = ContentId.deserialize(_mcid)
content_metadata = (await request.ctx.db_session.execute(select(StoredContent).where(StoredContent.hash == _cid.content_hash_b58))).scalars().first()
with open(content_metadata.filepath, 'r') as f:
content_metadata_json = json.loads(f.read())
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({
**opts,
'encrypted': content['encrypted_content'].json_format(),
'display_options': display_options
'display_options': display_options,
})
@@ -182,14 +219,17 @@ async def s_api_v1_content_friendly_list(request):
</tr>
</thead>
"""
for content in request.ctx.db_session.query(StoredContent).filter(
contents = (await request.ctx.db_session.execute(select(StoredContent).where(
StoredContent.type == 'onchain/content'
).all():
))).scalars().all()
for content in contents:
if not content.meta.get('metadata_cid'):
make_log("Content", f"Content {content.cid.serialize_v2()} has no metadata", level='warning')
continue
metadata_content = StoredContent.from_cid(request.ctx.db_session, content.meta.get('metadata_cid'))
from app.core.content.content_id import ContentId
_cid = ContentId.deserialize(content.meta.get('metadata_cid'))
metadata_content = (await request.ctx.db_session.execute(select(StoredContent).where(StoredContent.hash == _cid.content_hash_b58))).scalars().first()
with open(metadata_content.filepath, 'r') as f:
metadata = json.loads(f.read())
@@ -223,10 +263,12 @@ async def s_api_v1_5_content_list(request):
return response.json({'error': 'Invalid limit'}, status=400)
# Query onchain contents which are not disabled
contents = request.ctx.db_session.query(StoredContent).filter(
StoredContent.type == 'onchain/content',
StoredContent.disabled == False
).order_by(StoredContent.created.desc()).offset(offset).limit(limit).all()
contents = (await request.ctx.db_session.execute(
select(StoredContent)
.where(StoredContent.type == 'onchain/content', StoredContent.disabled == False)
.order_by(StoredContent.created.desc())
.offset(offset).limit(limit)
)).scalars().all()
result = []
for content in contents:
@@ -235,7 +277,9 @@ async def s_api_v1_5_content_list(request):
if not metadata_cid:
continue # Skip if no metadata_cid is found
metadata_content = StoredContent.from_cid(request.ctx.db_session, metadata_cid)
from app.core.content.content_id import ContentId
_cid = ContentId.deserialize(metadata_cid)
metadata_content = (await request.ctx.db_session.execute(select(StoredContent).where(StoredContent.hash == _cid.content_hash_b58))).scalars().first()
try:
with open(metadata_content.filepath, 'r') as f:
metadata = json.load(f)
@@ -251,9 +295,9 @@ async def s_api_v1_5_content_list(request):
preview_link = None
converted_content = content.meta.get('converted_content')
if converted_content:
converted_content = request.ctx.db_session.query(StoredContent).filter(
converted_content = (await request.ctx.db_session.execute(select(StoredContent).where(
StoredContent.hash == converted_content['low_preview']
).first()
))).scalars().first()
preview_link = converted_content.web_url
if converted_content.filename.split('.')[-1] in ('mp4', 'mov'):
media_type = 'video'
+27 -13
View File
@@ -11,6 +11,7 @@ from sanic import response
import json
from app.core._config import UPLOADS_DIR
from sqlalchemy import select
from app.core._utils.resolve_content import resolve_content
from app.core.logger import make_log
from app.core.models.node_storage import StoredContent
@@ -52,7 +53,9 @@ async def s_api_v1_storage_post(request):
try:
file_hash_bin = hashlib.sha256(file_content).digest()
file_hash = b58encode(file_hash_bin).decode()
stored_content = request.ctx.db_session.query(StoredContent).filter(StoredContent.hash == file_hash).first()
stored_content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == file_hash)
)).scalars().first()
if stored_content:
stored_cid = stored_content.cid.serialize_v1()
stored_cid_v2 = stored_content.cid.serialize_v2()
@@ -80,7 +83,7 @@ async def s_api_v1_storage_post(request):
key_id=None,
)
request.ctx.db_session.add(new_content)
request.ctx.db_session.commit()
await request.ctx.db_session.commit()
file_path = os.path.join(UPLOADS_DIR, file_hash)
async with aiofiles.open(file_path, "wb") as file:
@@ -97,7 +100,7 @@ async def s_api_v1_storage_post(request):
"content_url": f"dmy://storage?cid={new_cid}",
})
except BaseException as e:
make_log("Storage", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Error: {e}" + '\n' + traceback.format_exc(), level="error")
return response.json({"error": f"Error: {e}"}, status=500)
@@ -112,14 +115,16 @@ async def s_api_v1_storage_get(request, file_hash=None):
return response.json({"error": errmsg}, status=400)
content_sha256 = b58encode(cid.content_hash).decode()
content = request.ctx.db_session.query(StoredContent).filter(StoredContent.hash == content_sha256).first()
content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == content_sha256)
)).scalars().first()
if not content:
return response.json({"error": "File not found"}, status=404)
make_log("Storage", f"File {content_sha256} requested by {request.ctx.user}")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} File {content_sha256} requested by user={getattr(getattr(request.ctx, 'user', None), 'id', None)}")
file_path = os.path.join(UPLOADS_DIR, content_sha256)
if not os.path.exists(file_path):
make_log("Storage", f"File {content_sha256} not found locally", level="error")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} File {content_sha256} not found locally", level="error")
return response.json({"error": "File not found"}, status=404)
async with aiofiles.open(file_path, "rb") as file:
@@ -139,7 +144,16 @@ async def s_api_v1_storage_get(request, file_hash=None):
tempfile_path += "_mpeg" + (f"_{seconds_limit}" if seconds_limit else "")
if not os.path.exists(tempfile_path):
try:
cover_content = StoredContent.from_cid(content.meta.get('cover_cid'))
# Resolve cover content by CID (async)
from app.core.content.content_id import ContentId
try:
_cid = ContentId.deserialize(content.meta.get('cover_cid'))
_cover_hash = _cid.content_hash_b58
cover_content = (await request.ctx.db_session.execute(
select(StoredContent).where(StoredContent.hash == _cover_hash)
)).scalars().first()
except Exception:
cover_content = None
cover_tempfile_path = os.path.join(UPLOADS_DIR, f"tmp_{cover_content.hash}_jpeg")
if not os.path.exists(cover_tempfile_path):
cover_image = Image.open(cover_content.filepath)
@@ -173,25 +187,25 @@ async def s_api_v1_storage_get(request, file_hash=None):
try:
audio = AudioSegment.from_file(file_path)
except BaseException as e:
make_log("Storage", f"Error loading audio from file: {e}", level="debug")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Error loading audio from file: {e}", level="debug")
if not audio:
try:
audio = AudioSegment(content_file_bin)
except BaseException as e:
make_log("Storage", f"Error loading audio from binary: {e}", level="debug")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Error loading audio from binary: {e}", level="debug")
audio = audio[:seconds_limit * 1000] if seconds_limit else audio
audio.export(tempfile_path, format="mp3", cover=cover_tempfile_path)
except BaseException as e:
make_log("Storage", f"Error converting audio: {e}" + '\n' + traceback.format_exc(), level="error")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Error converting audio: {e}" + '\n' + traceback.format_exc(), level="error")
if os.path.exists(tempfile_path):
async with aiofiles.open(tempfile_path, "rb") as file:
content_file_bin = await file.read()
accept_type = 'audio/mpeg'
make_log("Storage", f"Audio {content_sha256} converted successfully")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Audio {content_sha256} converted successfully", level='debug')
else:
tempfile_path = tempfile_path[:-5]
@@ -208,13 +222,13 @@ async def s_api_v1_storage_get(request, file_hash=None):
break
quality -= 5
except BaseException as e:
make_log("Storage", f"Error converting image: {e}" + '\n' + traceback.format_exc(), level="error")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Error converting image: {e}" + '\n' + traceback.format_exc(), level="error")
if os.path.exists(tempfile_path):
async with aiofiles.open(tempfile_path, "rb") as file:
content_file_bin = await file.read()
make_log("Storage", f"Image {content_sha256} converted successfully")
make_log("Storage", f"sid={getattr(request.ctx, 'session_id', None)} Image {content_sha256} converted successfully", level='debug')
accept_type = 'image/jpeg'
else:
tempfile_path = tempfile_path[:-5]
+27 -26
View File
@@ -11,6 +11,7 @@ from base58 import b58encode
from sanic import response
from app.core.logger import make_log
from sqlalchemy import select
from app.core.models.node_storage import StoredContent
from app.core._config import UPLOADS_DIR
from app.core._utils.resolve_content import resolve_content
@@ -19,28 +20,28 @@ from app.core._utils.resolve_content import resolve_content
# POST /api/v1.5/storage
async def s_api_v1_5_storage_post(request):
# Log the receipt of a chunk upload request
make_log("uploader_v1.5", "Received chunk upload request", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Received chunk upload request", level="INFO")
# Get the provided file name from header and decode it from base64
provided_filename_b64 = request.headers.get("X-File-Name")
if not provided_filename_b64:
make_log("uploader_v1.5", "Missing X-File-Name header", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Missing X-File-Name header", level="ERROR")
return response.json({"error": "Missing X-File-Name header"}, status=400)
try:
provided_filename = b64decode(provided_filename_b64).decode("utf-8")
except Exception as e:
make_log("uploader_v1.5", f"Invalid X-File-Name header: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Invalid X-File-Name header: {e}", level="ERROR")
return response.json({"error": "Invalid X-File-Name header"}, status=400)
# Get X-Chunk-Start header (must be provided) and parse it as integer
chunk_start_header = request.headers.get("X-Chunk-Start")
if chunk_start_header is None:
make_log("uploader_v1.5", "Missing X-Chunk-Start header", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Missing X-Chunk-Start header", level="ERROR")
return response.json({"error": "Missing X-Chunk-Start header"}, status=400)
try:
chunk_start = int(chunk_start_header)
except Exception as e:
make_log("uploader_v1.5", f"Invalid X-Chunk-Start header: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Invalid X-Chunk-Start header: {e}", level="ERROR")
return response.json({"error": "Invalid X-Chunk-Start header"}, status=400)
# Enforce maximum chunk size (80 MB) using Content-Length header if provided
@@ -50,7 +51,7 @@ async def s_api_v1_5_storage_post(request):
try:
content_length = int(content_length)
if content_length > max_chunk_size:
make_log("uploader_v1.5", f"Chunk size {content_length} exceeds maximum allowed", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Chunk size {content_length} exceeds maximum allowed", level="ERROR")
return response.json({"error": "Chunk size exceeds maximum allowed (80 MB)"}, status=400)
except:
pass
@@ -62,9 +63,9 @@ async def s_api_v1_5_storage_post(request):
# New upload session: generate a new uuid
upload_id = str(uuid4())
is_new_upload = True
make_log("uploader_v1.5", f"Starting new upload session with ID: {upload_id}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Start new upload session id={upload_id}", level="INFO")
else:
make_log("uploader_v1.5", f"Resuming upload session with ID: {upload_id}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Resume upload session id={upload_id}", level="DEBUG")
# Determine the temporary file path based on upload_id
temp_path = os.path.join(UPLOADS_DIR, f"v1.5_upload_{upload_id}")
@@ -76,10 +77,10 @@ async def s_api_v1_5_storage_post(request):
# If the provided chunk_start is less than current_size, the chunk is already received
if chunk_start < current_size:
make_log("uploader_v1.5", f"Chunk starting at {chunk_start} already received, current size: {current_size}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Chunk at {chunk_start} already received; size={current_size}", level="DEBUG")
return response.json({"upload_id": upload_id, "current_size": current_size})
elif chunk_start > current_size:
make_log("uploader_v1.5", f"Chunk start {chunk_start} does not match current file size {current_size}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Chunk start {chunk_start} != current size {current_size}", level="ERROR")
return response.json({"error": "Chunk start does not match current file size"}, status=400)
# Append the received chunk to the temporary file
@@ -93,9 +94,9 @@ async def s_api_v1_5_storage_post(request):
async for chunk in request.stream:
await out_file.write(chunk)
new_size = os.path.getsize(temp_path)
make_log("uploader_v1.5", f"Appended chunk. New file size: {new_size}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Appended chunk. size={new_size}", level="DEBUG")
except Exception as e:
make_log("uploader_v1.5", f"Error saving chunk: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Error saving chunk: {e}", level="ERROR")
return response.json({"error": "Failed to save chunk"}, status=500)
# If computed hash matches the provided one, the final chunk has been received
@@ -111,28 +112,28 @@ async def s_api_v1_5_storage_post(request):
stdout, stderr = await proc.communicate()
if proc.returncode != 0:
error_msg = stderr.decode().strip()
make_log("uploader_v1.5", f"sha256sum error: {error_msg}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} sha256sum error: {error_msg}", level="ERROR")
return response.json({"error": "Failed to compute file hash"}, status=500)
computed_hash_hex = stdout.decode().split()[0].strip()
computed_hash_bytes = bytes.fromhex(computed_hash_hex)
computed_hash_b58 = b58encode(computed_hash_bytes).decode()
make_log("uploader_v1.5", f"Computed hash (base58): {computed_hash_b58}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Computed hash (base58): {computed_hash_b58}", level="INFO")
except Exception as e:
make_log("uploader_v1.5", f"Error computing file hash: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Error computing file hash: {e}", level="ERROR")
return response.json({"error": "Error computing file hash"}, status=500)
final_path = os.path.join(UPLOADS_DIR, f"{computed_hash_b58}")
try:
os.rename(temp_path, final_path)
make_log("uploader_v1.5", f"Final chunk received. File renamed to: {final_path}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Final chunk received. Renamed to: {final_path}", level="INFO")
except Exception as e:
make_log("uploader_v1.5", f"Error renaming file: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Error renaming file: {e}", level="ERROR")
return response.json({"error": "Failed to finalize file storage"}, status=500)
db_session = request.ctx.db_session
existing = db_session.query(StoredContent).filter_by(hash=computed_hash_b58).first()
existing = (await db_session.execute(select(StoredContent).where(StoredContent.hash == computed_hash_b58))).scalars().first()
if existing:
make_log("uploader_v1.5", f"File with hash {computed_hash_b58} already exists in DB", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} File already exists in DB: {computed_hash_b58}", level="INFO")
serialized_v2 = existing.cid.serialize_v2()
serialized_v1 = existing.cid.serialize_v1()
return response.json({
@@ -156,10 +157,10 @@ async def s_api_v1_5_storage_post(request):
created=datetime.utcnow()
)
db_session.add(new_content)
db_session.commit()
make_log("uploader_v1.5", f"New file stored and indexed for user {user_id} with hash {computed_hash_b58}", level="INFO")
await db_session.commit()
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Stored new file user={user_id} hash={computed_hash_b58}", level="INFO")
except Exception as e:
make_log("uploader_v1.5", f"Database error: {e}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Database error: {e}", level="ERROR")
return response.json({"error": "Database error"}, status=500)
serialized_v2 = new_content.cid.serialize_v2()
@@ -178,7 +179,7 @@ async def s_api_v1_5_storage_post(request):
# GET /api/v1.5/storage/<file_hash>
async def s_api_v1_5_storage_get(request, file_hash):
make_log("uploader_v1.5", f"Received file retrieval request for hash: {file_hash}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Retrieve file hash={file_hash}", level="INFO")
try:
file_hash = b58encode(resolve_content(file_hash)[0].content_hash).decode()
@@ -187,11 +188,11 @@ async def s_api_v1_5_storage_get(request, file_hash):
final_path = os.path.join(UPLOADS_DIR, f"{file_hash}")
if not os.path.exists(final_path):
make_log("uploader_v1.5", f"File not found: {final_path}", level="ERROR")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} File not found: {final_path}", level="ERROR")
return response.json({"error": "File not found"}, status=404)
db_session = request.ctx.db_session
stored = db_session.query(StoredContent).filter_by(hash=file_hash).first()
stored = (await db_session.execute(select(StoredContent).where(StoredContent.hash == file_hash))).scalars().first()
if stored and stored.filename:
filename_for_mime = stored.filename
else:
@@ -205,7 +206,7 @@ async def s_api_v1_5_storage_get(request, file_hash):
range_header = request.headers.get("Range")
if range_header:
make_log("uploader_v1.5", f"Processing Range header: {range_header}", level="INFO")
make_log("uploader_v1.5", f"sid={getattr(request.ctx, 'session_id', None)} Processing Range: {range_header}", level="DEBUG")
range_spec = range_header.strip().lower()
if not range_spec.startswith("bytes="):
make_log("uploader_v1.5", f"Invalid Range header: {range_header}", level="ERROR")
+18 -8
View File
@@ -4,6 +4,7 @@ from aiogram.utils.web_app import safe_parse_webapp_init_data
from sanic import response
from app.core._blockchain.ton.connect import TonConnect, unpack_wallet_info, WalletConnection
from sqlalchemy import select, and_
from app.core._config import TELEGRAM_API_KEY
from app.core.models.user import User
from app.core.logger import make_log
@@ -23,8 +24,19 @@ async def s_api_v1_tonconnect_new(request):
db_session = request.ctx.db_session
user = request.ctx.user
memory = request.ctx.memory
ton_connect, ton_connection = TonConnect.by_user(db_session, user)
# Try restore last connection from DB
ton_connection = (await db_session.execute(select(WalletConnection).where(
and_(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False,
WalletConnection.network == 'ton'
)
).order_by(WalletConnection.created.desc()))).scalars().first()
if ton_connection:
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"])
await ton_connect.restore_connection()
else:
ton_connect = TonConnect()
make_log("TonConnect_API", f"SDK connected?: {ton_connect.connected}", level='info')
if ton_connect.connected:
return response.json({"error": "Already connected"}, status=400)
@@ -47,13 +59,11 @@ async def s_api_v1_tonconnect_logout(request):
user = request.ctx.user
memory = request.ctx.memory
wallet_connections = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).all()
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))
wallet_connections = result.scalars().all()
for wallet_connection in wallet_connections:
wallet_connection.invalidated = True
db_session.commit()
await db_session.commit()
return response.json({"success": True})
+11 -10
View File
@@ -1,6 +1,7 @@
from app.core.logger import make_log, logger
from app.core.models._telegram import Wrapped_CBotChat
from app.core.models.user import User
from sqlalchemy import select
from app.core.storage import db_session
from aiogram import BaseMiddleware, types
from app.core.models.messages import KnownTelegramMessage
@@ -21,9 +22,9 @@ class UserDataMiddleware(BaseMiddleware):
# TODO: maybe make users cache
with db_session(auto_commit=False) as session:
async with db_session(auto_commit=False) as session:
try:
user = session.query(User).filter_by(telegram_id=user_id).first()
user = (await session.execute(select(User).where(User.telegram_id == user_id))).scalars().first()
except BaseException as e:
logger.error(f"Error when middleware getting user: {e}")
user = None
@@ -42,7 +43,7 @@ class UserDataMiddleware(BaseMiddleware):
created=datetime.now()
)
session.add(user)
session.commit()
await session.commit()
else:
if user.username != update_body.from_user.username:
user.username = update_body.from_user.username
@@ -60,7 +61,7 @@ class UserDataMiddleware(BaseMiddleware):
}
user.last_use = datetime.now()
session.commit()
await session.commit()
data['user'] = user
data['db_session'] = session
@@ -72,11 +73,11 @@ class UserDataMiddleware(BaseMiddleware):
if update_body.text.startswith('/start'):
message_type = 'start_command'
if session.query(KnownTelegramMessage).filter_by(
chat_id=update_body.chat.id,
message_id=update_body.message_id,
from_user=True
).first():
if (await session.execute(select(KnownTelegramMessage).where(
(KnownTelegramMessage.chat_id == update_body.chat.id) &
(KnownTelegramMessage.message_id == update_body.message_id) &
(KnownTelegramMessage.from_user == True)
))).scalars().first():
make_log("UserDataMiddleware", f"Message {update_body.message_id} already processed", level='debug')
return
@@ -91,7 +92,7 @@ class UserDataMiddleware(BaseMiddleware):
meta={}
)
session.add(new_message)
session.commit()
await session.commit()
result = await handler(event, data)
return result
+9 -8
View File
@@ -6,6 +6,7 @@ from app.core._keyboards import get_inline_keyboard
from app.core._utils.tg_process_template import tg_process_template
from app.core.logger import make_log
from app.core.models.node_storage import StoredContent
from sqlalchemy import select, and_
import json
router = Router()
@@ -20,12 +21,13 @@ def chunks(lst, n):
async def t_callback_owned_content(query: types.CallbackQuery, memory=None, user=None, db_session=None, chat_wrap=None, **extra):
message_text = user.translated("ownedContent_menu")
content_list = []
for content in db_session.query(StoredContent).filter_by(
owner_address=user.wallet_address(db_session),
type='onchain/content'
).all():
user_addr = await user.wallet_address_async(db_session)
result = await db_session.execute(select(StoredContent).where(
and_(StoredContent.owner_address == user_addr, StoredContent.type == 'onchain/content')
))
for content in result.scalars().all():
try:
metadata_content = StoredContent.from_cid(db_session, content.json_format()['metadata_cid'])
metadata_content = await StoredContent.from_cid_async(db_session, content.json_format()['metadata_cid'])
with open(metadata_content.filepath, 'r') as f:
metadata_content_json = json.loads(f.read())
except BaseException as e:
@@ -59,10 +61,9 @@ async def t_callback_owned_content(query: types.CallbackQuery, memory=None, user
async def t_callback_node_content(query: types.CallbackQuery, memory=None, user=None, db_session=None, chat_wrap=None, **extra):
content_oid = int(query.data.split('_')[1])
row = (await db_session.execute(select(StoredContent).where(StoredContent.id == content_oid))).scalars().first()
return await chat_wrap.send_content(
db_session, db_session.query(StoredContent).filter_by(
id=content_oid
).first(),
db_session, row,
extra_buttons=[
[{
'text': user.translated('back_button'),
+11 -5
View File
@@ -3,6 +3,7 @@ from aiogram.filters import Command
from tonsdk.utils import Address
from app.core._blockchain.ton.connect import TonConnect
from sqlalchemy import select, and_
from app.core._keyboards import get_inline_keyboard
from app.core._utils.tg_process_template import tg_process_template
from app.core.models.wallet_connection import WalletConnection
@@ -32,7 +33,13 @@ async def send_home_menu(chat_wrap, user, wallet_connection, **kwargs):
async def send_connect_wallets_list(db_session, chat_wrap, user, **kwargs):
ton_connect, ton_connection = TonConnect.by_user(db_session, user, callback_fn=())
# Try to restore existing connection via DB
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
if ton_connection:
await ton_connect.restore_connection()
wallets = ton_connect._sdk_client.get_wallets()
message_text = user.translated("connectWalletsList_menu")
@@ -66,10 +73,9 @@ async def t_home_menu(__msg, **extra):
else:
message_id = None
wallet_connection = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).first()
wallet_connection = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().first()
# if not wallet_connection:
# return await send_connect_wallets_list(db_session, chat_wrap, user, message_id=message_id)
+22 -12
View File
@@ -7,6 +7,7 @@ from aiogram.filters import Command
from app.bot.routers.home import send_connect_wallets_list, send_home_menu
from app.core._blockchain.ton.connect import TonConnect, unpack_wallet_info
from sqlalchemy import select, and_
from app.core._keyboards import get_inline_keyboard
from app.core._utils.tg_process_template import tg_process_template
from app.core.logger import make_log
@@ -33,15 +34,21 @@ async def t_tonconnect_dev_menu(message: types.Message, memory=None, user=None,
keyboard = []
ton_connect, ton_connection = TonConnect.by_user(db_session, user, callback_fn=())
# Restore recent connection
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
make_log("TonConnect_DevMenu", f"Available wallets: {ton_connect._sdk_client.get_wallets()}", level='debug')
if ton_connection:
await ton_connect.restore_connection()
make_log("TonConnect_DevMenu", f"SDK connected?: {ton_connect.connected}", level='info')
if not ton_connect.connected:
if ton_connection:
make_log("TonConnect_DevMenu", f"Invalidating old connection", level='debug')
ton_connection.invalidated = True
db_session.commit()
await db_session.commit()
message_text = f"""<b>Wallet is not connected</b>
@@ -71,7 +78,12 @@ Use /dev_tonconnect <code>{wallet_app_name}</code> for connect to wallet."""
async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, user=None, db_session=None, chat_wrap=None, **extra):
wallet_app_name = query.data.split("_")[1]
ton_connect, ton_connection = TonConnect.by_user(db_session, user)
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
if ton_connection:
await ton_connect.restore_connection()
connection_link = await ton_connect.new_connection(wallet_app_name)
ton_connect.connected
@@ -98,10 +110,9 @@ async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, us
start_ts = datetime.now()
while datetime.now() - start_ts < timedelta(seconds=180):
new_connection = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).first()
new_connection = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().first()
if new_connection:
await tg_process_template(
chat_wrap, user.translated('p_successConnectWallet')
@@ -115,14 +126,13 @@ async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, us
async def t_callback_disconnect_wallet(query: types.CallbackQuery, memory=None, user=None, db_session=None, chat_wrap=None, **extra):
wallet_connections = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).all()
wallet_connections = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().all()
for wallet_connection in wallet_connections:
wallet_connection.invalidated = True
db_session.commit()
await db_session.commit()
return await send_home_menu(chat_wrap, user, None, message_id=query.message.message_id)
+23 -24
View File
@@ -6,6 +6,7 @@ from aiogram import types, Router, F
from app.core._keyboards import get_inline_keyboard
from app.core.models.node_storage import StoredContent
from sqlalchemy import select, and_
import json
from app.core.logger import make_log
from app.core.models.content.user_content import UserAction, UserContent
@@ -30,7 +31,7 @@ CACHE_CHAT_ID = -1002390124789
async def t_callback_purchase_node_content(query: types.CallbackQuery, memory=None, user=None, db_session=None, chat_wrap=None, **extra):
content_oid = int(query.data.split('_')[1])
is_cancel_request = query.data.split('_')[2] == 'cancel' if len(query.data.split('_')) > 2 else False
content = db_session.query(StoredContent).filter_by(id=content_oid).first()
content = (await db_session.execute(select(StoredContent).where(StoredContent.id == content_oid))).scalars().first()
if not content:
return await query.answer(user.translated('error_contentNotFound'), show_alert=True)
@@ -43,11 +44,16 @@ async def t_callback_purchase_node_content(query: types.CallbackQuery, memory=No
make_log("Purchase", f"User {user.id} initiated purchase for content ID {content_oid}. License price: {license_price_num}.", level='info')
ton_connect, ton_connection = TonConnect.by_user(db_session, user, callback_fn=())
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
if ton_connection:
await ton_connect.restore_connection()
assert ton_connect.connected, "No connected wallet"
user_wallet_address = user.wallet_address(db_session)
user_wallet_address = await user.wallet_address_async(db_session)
memory._app.add_task(ton_connect._sdk_client.send_transaction({
'valid_until': int(datetime.now().timestamp() + 300),
@@ -76,18 +82,15 @@ async def t_callback_purchase_node_content(query: types.CallbackQuery, memory=No
else:
# Logging cancellation attempt with detailed information
make_log("Purchase", f"User {user.id} cancelled purchase for content ID {content_oid}.", level='info')
action = db_session.query(UserAction).filter_by(
type='purchase',
content_id=content_oid,
user_id=user.id,
status='requested'
).first()
action = (await db_session.execute(select(UserAction).where(
and_(UserAction.type == 'purchase', UserAction.content_id == content_oid, UserAction.user_id == user.id, UserAction.status == 'requested')
))).scalars().first()
if not action:
return await query.answer()
action.status = 'canceled'
db_session.commit()
await db_session.commit()
await chat_wrap.send_content(db_session, content, message_id=query.message.message_id)
@@ -104,9 +107,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
args = None
if source_args_ext.startswith('Q'):
license_onchain_address = source_args_ext[1:]
licensed_content = db_session.query(UserContent).filter_by(
onchain_address=license_onchain_address,
).first().content
licensed_content = (await db_session.execute(select(UserContent).where(UserContent.onchain_address == license_onchain_address))).scalars().first().content
make_log("InlineSearch", f"Query '{query.query}' is a license query for content ID {licensed_content.id}.", level='info')
args = licensed_content.cid.serialize_v2()
else:
@@ -118,15 +119,15 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
content_list = []
search_query = {'hash': cid.content_hash_b58}
make_log("InlineSearch", f"Searching with query '{search_query}'.", level='info')
content = db_session.query(StoredContent).filter_by(**search_query).first()
content_prod = content.open_content(db_session)
content = (await db_session.execute(select(StoredContent).where(StoredContent.hash == cid.content_hash_b58))).scalars().first()
content_prod = await content.open_content_async(db_session)
# Get both encrypted and decrypted content objects
encrypted_content = content_prod['encrypted_content']
decrypted_content = content_prod['decrypted_content']
decrypted_content_meta = decrypted_content.json_format()
try:
metadata_content = StoredContent.from_cid(db_session, content.json_format()['metadata_cid'])
metadata_content = await StoredContent.from_cid_async(db_session, content.json_format()['metadata_cid'])
with open(metadata_content.filepath, 'r') as f:
metadata_content_json = json.loads(f.read())
except BaseException as e:
@@ -144,7 +145,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
result_kwargs = {}
try:
cover_content = StoredContent.from_cid(db_session, decrypted_content_meta.get('cover_cid') or None)
cover_content = await StoredContent.from_cid_async(db_session, decrypted_content_meta.get('cover_cid') or None)
except BaseException as e:
cover_content = None
@@ -152,9 +153,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
result_kwargs['thumb_url'] = cover_content.web_url
content_type_declared = decrypted_content_meta.get('content_type', 'application/x-binary').split('/')[0]
preview_content = db_session.query(StoredContent).filter_by(
hash=content.meta.get('converted_content', {}).get('low_preview')
).first()
preview_content = (await db_session.execute(select(StoredContent).where(StoredContent.hash == content.meta.get('converted_content', {}).get('low_preview')))).scalars().first()
content_type_declared = {
'mp3': 'audio',
'flac': 'audio',
@@ -162,7 +161,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
'mov': 'video'
}.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:
hashtags_str = hashtags_str + '\n'
@@ -196,7 +195,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
**decrypted_content.meta,
'telegram_file_cache_preview': preview_file_id
}
db_session.commit()
await db_session.commit()
except Exception as e:
# Logging error during preview upload with detailed content type and query information
make_log("InlineSearch", f"Error uploading preview for content type '{content_type_declared}' during inline query '{query.query}': {e}", level='error')
@@ -212,7 +211,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
types.InlineQueryResultCachedAudio(
id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}",
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',
reply_markup=get_inline_keyboard([
[
@@ -242,7 +241,7 @@ async def t_inline_query_node_content(query: types.InlineQuery, memory=None, use
id=f"NC_{content.id}_{int(datetime.now().timestamp() // 60)}",
video_file_id=decrypted_content.meta['telegram_file_cache_preview'],
title=title,
caption=hashtags_str + user.translated('p_playerContext_preview'),
caption=hashtags_str + user.translated('p_playerVideoContext_preview'),
parse_mode='html',
reply_markup=get_inline_keyboard([
[
+10 -5
View File
@@ -3,6 +3,7 @@ from aiogram.filters import Command
from tonsdk.utils import Address
from app.core._blockchain.ton.connect import TonConnect
from sqlalchemy import select, and_
from app.core._keyboards import get_inline_keyboard
from app.core._utils.tg_process_template import tg_process_template
from app.core.logger import make_log
@@ -32,7 +33,12 @@ async def send_home_menu(chat_wrap, user, wallet_connection, **kwargs):
async def send_connect_wallets_list(db_session, chat_wrap, user, **kwargs):
ton_connect, ton_connection = TonConnect.by_user(db_session, user, callback_fn=())
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
if ton_connection:
await ton_connect.restore_connection()
wallets = ton_connect._sdk_client.get_wallets()
message_text = user.translated("connectWalletsList_menu")
@@ -66,10 +72,9 @@ async def t_home_menu(__msg, **extra):
else:
message_id = None
wallet_connection = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).first()
wallet_connection = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().first()
# if not wallet_connection:
# return await send_connect_wallets_list(db_session, chat_wrap, user, message_id=message_id)
+21 -12
View File
@@ -7,6 +7,7 @@ from aiogram.filters import Command
from app.client_bot.routers.home import send_connect_wallets_list, send_home_menu
from app.core._blockchain.ton.connect import TonConnect, unpack_wallet_info
from sqlalchemy import select, and_
from app.core._keyboards import get_inline_keyboard
from app.core._utils.tg_process_template import tg_process_template
from app.core.logger import make_log
@@ -34,15 +35,20 @@ async def t_tonconnect_dev_menu(message: types.Message, memory=None, user=None,
keyboard = []
ton_connect, ton_connection = TonConnect.by_user(db_session, user, callback_fn=())
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
make_log("TonConnect_DevMenu", f"Available wallets: {ton_connect._sdk_client.get_wallets()}", level='debug')
if ton_connection:
await ton_connect.restore_connection()
make_log("TonConnect_DevMenu", f"SDK connected?: {ton_connect.connected}", level='info')
if not ton_connect.connected:
if ton_connection:
make_log("TonConnect_DevMenu", f"Invalidating old connection", level='debug')
ton_connection.invalidated = True
db_session.commit()
await db_session.commit()
message_text = f"""<b>Wallet is not connected</b>
@@ -73,7 +79,12 @@ Use /dev_tonconnect <code>{wallet_app_name}</code> for connect to wallet."""
async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, user=None, db_session=None,
chat_wrap=None, **extra):
wallet_app_name = query.data.split("_")[1]
ton_connect, ton_connection = TonConnect.by_user(db_session, user)
result = await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False, WalletConnection.network == 'ton')
).order_by(WalletConnection.created.desc()))
ton_connection = result.scalars().first()
ton_connect = TonConnect.by_key(ton_connection.keys["connection_key"]) if ton_connection else TonConnect()
if ton_connection:
await ton_connect.restore_connection()
connection_link = await ton_connect.new_connection(wallet_app_name)
ton_connect.connected
@@ -100,10 +111,9 @@ async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, us
start_ts = datetime.now()
while datetime.now() - start_ts < timedelta(seconds=180):
new_connection = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).first()
new_connection = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().first()
if new_connection:
await tg_process_template(
chat_wrap, user.translated('p_successConnectWallet')
@@ -118,14 +128,13 @@ async def t_callback_init_tonconnect(query: types.CallbackQuery, memory=None, us
async def t_callback_disconnect_wallet(query: types.CallbackQuery, memory=None, user=None, db_session=None,
chat_wrap=None, **extra):
wallet_connections = db_session.query(WalletConnection).filter(
WalletConnection.user_id == user.id,
WalletConnection.invalidated == False
).all()
wallet_connections = (await db_session.execute(select(WalletConnection).where(
and_(WalletConnection.user_id == user.id, WalletConnection.invalidated == False)
))).scalars().all()
for wallet_connection in wallet_connections:
wallet_connection.invalidated = True
db_session.commit()
await db_session.commit()
return await send_home_menu(chat_wrap, user, None, message_id=query.message.message_id)
+1 -1
View File
@@ -5,7 +5,7 @@ from app.core._secrets import service_wallet
class Blank(Contract):
code = 'B5EE9C72010104010042000114FF00F4A413F4BCF2C80B010202CA0203004FD043A0E9AE43F48061DA89A1F480618E0BE5C323A803A1A843F60803A1DA3DDAA7A861DAA9E2026F0007A0DD7C12'
code = 'b5ee9c72410104010042000114ff00f4a413f4bcf2c80b010202ca03020007a0dd7c12004fd043a0e9ae43f48061da89a1f480618e0be5c323a803a1a843f60803a1da3ddaa7a861daa9e2026f102bdd33'
def __init__(self, **kwargs):
kwargs['code'] = Cell.one_from_boc(self.code)
@@ -6,7 +6,7 @@ from app.core._config import MY_PLATFORM_CONTRACT
class COP_NFT(Contract):
code = 'b5ee9c7241022f01000ac8000114ff00f4a413f4bcf2c80b010201620b020201200803020120070402016a0605004bae43b941781bb12d2f0dfa54dda7472bd6fbc06813c0f71615c0094ee3548791d80b2bfcecc0007dac1ff6a26869ff80fc317d2000fc31906ba4e1807c30fc20c70d698380fc327d2000fc32ea00fc336a00fc33ea00fc346a00fc34ef68fc21fc217c227c22c000cbb8fcfed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded1f8419ac8c9707f22d023d05503e1f849d0f846f00702987ff84504d4305e21e05b7ff842f843f84504d4301034413080201200a0900d5ba7a3ed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded1f848d0d431d430f8287003c8cbffc94130c85003cf16cb07ccc97020c8cb0113f400f400cb00c9f9007074c8cb02ca07cbffc9d080093bb54ded44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded171f828f842f843f844f845f849f846f848f84780202c71e0c020148180d020272140e0101f40f01bac882f01bc04b5291c26a46d918139138b992d2de976d6851d0893b0476b85bfbdfc6e6228307f40e6fa130d3ff3001cbff82f03dbe3c57aae062079b4712b8a650cf2abe2d6f986e0ecf1785d2ba835974175d228307f40e6fa130cf1610015e82f0e3868b37bc801236c714f134200a4e3211860f3ad67d22be995321452723c516228307f40e6fa1923031e30dc91101c0d3073001cb0782f0862bd78e42e143bb51660c521752437648ef67967d5055093d11be1117c774b0228307f40e6fa130cf1682f0ade34d680d2df2ad88267395ea1c3b159c4df766416d69d2bbcc8e18e87bd2ba228307f40e6fa130d43001cc1201fe82f093354845030274cd4bf1686abd60ab28ec52e1a792fa6f7ccb9cbd0ddff53d12228307f40e6fa130d43001cc82f0852d0f3fd8bdef2aab5f91714d86638ddc57db97743a350fab3394bc9ebd9518228307f40e6fa130d43001cc82f089445ea08b55421faa49919a5fd272e9a520f701b479d6084847e161ca5b7711581300168307f40e6fa130d43001cc01bbf36fc216465ffc1780de025a948e135236c8c09c89c5cc9696f4bb6b428e8449d823b5c2dfdefe3732c4183fa21e47c21e78b41781edf1e2bd5703103cda3895c532867955f16b7cc3707678bc2e95d41acba0baeac4183fa21fc20f18041501fef844c8cb0782f0e3868b37bc801236c714f134200a4e3211860f3ad67d22be995321452723c516588307f443c8f845cf1682f0862bd78e42e143bb51660c521752437648ef67967d5055093d11be1117c774b0588307f443f846c8cc82f0ade34d680d2df2ad88267395ea1c3b159c4df766416d69d2bbcc8e18e87bd2ba581601bc8307f443f847c8cc82f093354845030274cd4bf1686abd60ab28ec52e1a792fa6f7ccb9cbd0ddff53d12588307f443f848c8cc82f0852d0f3fd8bdef2aab5f91714d86638ddc57db97743a350fab3394bc9ebd9518588307f443f849c8cc17004e82f089445ea08b55421faa49919a5fd272e9a520f701b479d6084847e161ca5b7711588307f4430202761b19017df1998e83a6b90fd201876a26869ff80fc317d2000fc31906ba4e1807c30fc20c70d698380fc327d2000fc32ea00fc336a00fc33ea00fc346a00fc34ef6880c1a01c8d3fffa4001f865f8435230c7058e5331fa405312c7058e416c12d30701f864f84271c8cb00cb3f58cf16c9f866d401f867d401f868d430f869f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54e05bf843c705f2e1afe30d2e0201201d1c00415f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54800694ed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded180201cb201f0011f686900699ffd20184020148222100113e910c30003cb8536003f5007434c0c05c6c2497c1383e903e900c7e800c5c75c87e800c7e800c1cea6d003b513434ffc07e18be90007e18c835d270c03e187e106386b4c1c07e193e90007e1975007e19b5007e19f5007e1a35007e1a77b47e1078c0dc14c0f5d270882616c0b4c7f4cfd111378860840a8c6564eea3a116c4b6cf380d48202d272302da82105fcc3d14ba8e85304330db3ce031342382102fcb26a2ba8eb23132f846f007706d82108b771735c8cb1f16cb3f049702c8cbff01cf169b6c21f842c8cbfff843cf16e212cf17128040db3ce0320282102fa30f96ba9ff843c705f2e191d401fb04d430ed54e05b840ff2f0242c03e0f84514c705f2e191f844c000f844c003b1f2e191fa4021f001fa40d20031fa000682084c4b40a121945315a0a1de22d70b01c300209206a19136e220c2fff2e19223f865218e9d6d821005138d91c8cb1f5260cb3ff845cf165008cf1610341771db3c039410266c31e202925f03e30d2c26250040f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54012622f0016d8210d53276dbc8cb1f12cb3f71db3c2c01e6f847d0d403d30721c003f2e19af844c000f844c003b1f2e1a4d3fffa403020d70b01c000923025de05f4043003d022a59320c2009601fa003101a5e830fa0030208209312d005320bcf2e208f848d0d4d430f8287007c8cbffc9225088c85003cf16cb07ccc97020c8cb0113f400f400cb00c92803fef9007074c8cb02ca07cbffc9d01ac705965392bef2e208df5192a1f846f007228f5c50675f053333343434028107d0a8812710a904f8456d70c8cb1f8d04935648149959995c9c985b0814185e5bdd5d20cf16102372db3c82084c4b4070fb02706d732282102a319593c8cb1fcb3fcb0715cbff5003cf1613810082db3ce02c2c2901fe35f8287022c8cbffc91038c85003cf16cb07ccc97020c8cb0113f400f400cb00c920f9007074c8cb02ca07cbffc9d0f849d0c803d013cf16f849f8486d6dc870fa0270fa02500afa02c9c8cc19f40018f400c904d3ff3029c8cbfff843cf16c90fc8cc1fccc9c8cc1ecbff2ccf16f828cf1619cb0712cc14cc1acc10234a502a03b472db3ca405c8ca0015cb3f5005cf16c9f866f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed5422c2008e926d708210d53276dbc8cb1fcb3f102472db3c926c21e22282103b9aca00bc925f03e30d2c2c2b02e202821005f5e100a1208106a4a8812710a90466a1f8436d8210247a9ebac8cb1f1023102472db3c593121c2008ec08ebc78f4966fa531208eaf01d430d0fa40d30f305240a8812710a9046d70c8cb1f8d04535648149Line truncated
code = 'b5ee9c7241023201000bbc000114ff00f4a413f4bcf2c80b01020120040201f6f2ed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded1f849d0d3ffd31f038308d71820f9015882f053b50ddfe5a9533f2e76ac054411db94432a1f7b7ae17fc64cf7aec5df8705d5f910f2e212fa40f82812c705f2e213d31f5213bc0301d4f2e21401d33f01f823bcf2e215f80082084c4b4070fb02d4308e3cd0d31f218210e3e30001ba8e1731f404216e91319301fb04e2f404216e91319301ed54e28e10018210e3e30002ba96d307d402fb00dee2f40430206ee63002d4d43002c8cbff13cb1f12ccccc9f869280201480e050201200b060201200a0702016a0908004bae43bac17829da86eff2d4a99f973b5602a208edca21950fbdbd70bfe3267bd762efc382eac0007dac1ff6a26869ff80fc317d2000fc31906ba4e1807c30fc20c70d698380fc327d2000fc32ea00fc336a00fc33ea00fc346a00fc34ef68fc21fc217c227c22c000cbb8fcfed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded1f8419ac8c9707f22d023d05503e1f849d0f846f00702987ff84504d4305e21e05b7ff842f843f84504d4301034413080201200d0c00d5ba7a3ed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded1f848d0d431d430f8287003c8cbffc94130c85003cf16cb07ccc97020c8cb0113f400f400cb00c9f9007074c8cb02ca07cbffc9d080093bb54ded44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded171f828f842f843f844f845f849f846f848f84780202c7210f0201481b1002027217110101f41201bac882f01bc04b5291c26a46d918139138b992d2de976d6851d0893b0476b85bfbdfc6e6228307f40e6fa130d3ff3001cbff82f03dbe3c57aae062079b4712b8a650cf2abe2d6f986e0ecf1785d2ba835974175d228307f40e6fa130cf1613015e82f0e3868b37bc801236c714f134200a4e3211860f3ad67d22be995321452723c516228307f40e6fa1923031e30dc91401c0d3073001cb0782f0862bd78e42e143bb51660c521752437648ef67967d5055093d11be1117c774b0228307f40e6fa130cf1682f0ade34d680d2df2ad88267395ea1c3b159c4df766416d69d2bbcc8e18e87bd2ba228307f40e6fa130d43001cc1501fe82f093354845030274cd4bf1686abd60ab28ec52e1a792fa6f7ccb9cbd0ddff53d12228307f40e6fa130d43001cc82f0852d0f3fd8bdef2aab5f91714d86638ddc57db97743a350fab3394bc9ebd9518228307f40e6fa130d43001cc82f089445ea08b55421faa49919a5fd272e9a520f701b479d6084847e161ca5b7711581600168307f40e6fa130d43001cc01bbf36fc216465ffc1780de025a948e135236c8c09c89c5cc9696f4bb6b428e8449d823b5c2dfdefe3732c4183fa21e47c21e78b41781edf1e2bd5703103cda3895c532867955f16b7cc3707678bc2e95d41acba0baeac4183fa21fc20f18041801fef844c8cb0782f0e3868b37bc801236c714f134200a4e3211860f3ad67d22be995321452723c516588307f443c8f845cf1682f0862bd78e42e143bb51660c521752437648ef67967d5055093d11be1117c774b0588307f443f846c8cc82f0ade34d680d2df2ad88267395ea1c3b159c4df766416d69d2bbcc8e18e87bd2ba581901bc8307f443f847c8cc82f093354845030274cd4bf1686abd60ab28ec52e1a792fa6f7ccb9cbd0ddff53d12588307f443f848c8cc82f0852d0f3fd8bdef2aab5f91714d86638ddc57db97743a350fab3394bc9ebd9518588307f443f849c8cc1a004e82f089445ea08b55421faa49919a5fd272e9a520f701b479d6084847e161ca5b7711588307f4430202761e1c017df1998e83a6b90fd201876a26869ff80fc317d2000fc31906ba4e1807c30fc20c70d698380fc327d2000fc32ea00fc336a00fc33ea00fc346a00fc34ef6880c1d01c8d3fffa4001f865f8435230c7058e5331fa405312c7058e416c12d30701f864f84271c8cb00cb3f58cf16c9f866d401f867d401f868d430f869f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54e05bf843c705f2e1afe30d31020120201f00415f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54800694ed44d0d3ff01f862fa4001f86320d749c300f861f8418e1ad30701f864fa4001f865d401f866d401f867d401f868d401f869ded180201cb23220011f686900699ffd20184020148252400113e910c30003cb8536003f5007434c0c05c6c2497c1383e903e900c7e800c5c75c87e800c7e800c1cea6d003b513434ffc07e18be90007e18c835d270c03e187e106386b4c1c07e193e90007e1975007e19b5007e19f5007e1a35007e1a77b47e1078c0dc14c0f5d270882616c0b4c7f4cfd111378860840a8c6564eea3a116c4b6cf380d4820302a2602da82105fcc3d14ba8e85304330db3ce031342382102fcb26a2ba8eb23132f846f007706d82108b771735c8cb1f16cb3f049702c8cbff01cf169b6c21f842c8cbfff843cf16e212cf17128040db3ce0320282102fa30f96ba9ff843c705f2e191d401fb04d430ed54e05b840ff2f0272f03e0f84514c705f2e191f844c000f844c003b1f2e191fa4021f001fa40d20031fa000682084c4b40a121945315a0a1de22d70b01c300209206a19136e220c2fff2e19223f865218e9d6d821005138d91c8cb1f5260cb3ff845cf165008cf1610341771db3c039410266c31e202925f03e30d2f29280040f849f848f847f846f844f842c8cbfff843cf16cb07f845cf16ccccccccc9ed54012622f0016d8210d53276dbc8cb1f12cb3f71db3c2f01e6f847d0d403d30721c003f2e19af844c000f844c003b1f2e1a4d3fffa403020d70b01c000923025de05f4043003d022a59320c2009601fa003101a5e830fa0030208209312d005320bcf2e208f848d0d4d430f8287007c8cbffc9225088c85003cf16cb07ccc97020c8cb0113f400f400cb00c92b03fef9007074c8cb02ca07cbffc9d01ac705965392bef2e208df5192a1f846f007228f5c50675f053333343434028107d0a8812710a904f8456d70c8cb1f8d04935648149959995c9c985b0814185e5bdd5d20cf16102372db3c82084c4b4070fb02706d732282102a319593c8cb1fcb3fcb0715cbff5003cf1613810082db3ce02f2f2c01fe35f8287022c8cbffc91038c85003cf16cb07ccc97020c8cb0113f400f400cb00c920f9007074c8cb02ca07cbffc9d0f849d0c803d013cf16f849f8486d6dc870fa0270fa02500Line truncated
codebase_version = 5
def __init__(self, **kwargs):
@@ -3,7 +3,7 @@ from tonsdk.contract import Contract
class Platform(Contract):
code = 'b5ee9c7241021601000310000114ff00f4a413f4bcf2c80b010201620d0202012006030201200504004bbac877282f037625a5e1bf4a9bb4e8e57adf780d02781ee2c2b80129dc6a90f23b01657f9d980057b905bed44d0fa4001f861d3ff01f862d401f863f843d0d431d430f864d401f865d1f845d0f84201d430f84180201200a07020120090800a1b4f47da89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba3f089f050e0079197ff92826190a0079e2d960f9992e04191960227e801e801960193f200e0e9919605940f97ff93a10000fb5daeeb00c9f05100201200c0b0059b6a9bda89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba2e1f051f085f087f089f08b00051b56ba63da89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba2e391960f999300202c70f0e0007a0dd7c120201cf111000113e910c30003cb8536002f30cf434c0c05c6c2497c0f83e90087c007e900c7e800c5c75c87e800c7e800c1cea6d0008f5d27048245c2540f4c7d411388830002497c1783b51343e90007e1874ffc07e18b5007e18fe10f4350c750c3e1935007e1974482084091ea7aeaea497c178082084152474232ea3a14c104c36cf380c4cbe1071c160131201dcf2e19120820833cc77ba9730d4d30730fb00e0208210b99cd03bba9701fa4001f86101de208210d81c632fba9601d401f86501de208210b5de5f9eba8e8b30fa40fa00306d6d71db3ce082102fa30f96ba98d401fb04d430ed54e030f845f843f842c8f841cf16cbffccccc9ed541502f082084c4b4001a013bef2e20801d3ffd4d430f844f82870f842c8cbffc9c85003cf16cb07ccc97020c8cb0113f400f400cb00c920f9007074c8cb02ca07cbffc9d0f843d070c804d014cf16f843f842c8cbfff828cf16c903d430c8cc13ccc9c8cc17cbff5007cf1614cc15cccc43308040db3cf842a4f86215140024f845f843f842c8f841cf16cbffccccc9ed540078708010c8cb055006cf165004fa0214cb68216e947032cb019bc858cf17c97158cb00f400e2226e95327058cb0099c85003cf17c958f400e2c901fb004e32cb65'
code = 'b5ee9c724102160100032e000114ff00f4a413f4bcf2c80b010201620d0202012006030201200504004bbac877582f053b50ddfe5a9533f2e76ac054411db94432a1f7b7ae17fc64cf7aec5df8705d580057b905bed44d0fa4001f861d3ff01f862d401f863f843d0d431d430f864d401f865d1f845d0f84201d430f84180201200a07020120090800a1b4f47da89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba3f089f050e0079197ff92826190a0079e2d960f9992e04191960227e801e801960193f200e0e9919605940f97ff93a10000fb5daeeb00c9f05100201200c0b0059b6a9bda89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba2e1f051f085f087f089f08b00051b56ba63da89a1f48003f0c3a7fe03f0c5a803f0c7f087a1a863a861f0c9a803f0cba2e391960f999300202c70f0e0007a0dd7c120201cf111000113e910c30003cb8536002f30cf434c0c05c6c2497c0f83e90087c007e900c7e800c5c75c87e800c7e800c1cea6d0008f5d27048245c2540f4c7d411388830002497c1783b51343e90007e1874ffc07e18b5007e18fe10f4350c750c3e1935007e1974482084091ea7aeaea497c178082084152474232ea3a14c104c36cf380c4cbe1071c160131201faf2e19120820833cc77ba9730d4d30730fb00e0208210b99cd03bba9701fa4001f86101de208210d81c632fba9601d401f86501de208210b5de5f9eba8e8b30fa40fa00306d6d71db3ce082102fa30f96ba8e16f404216e91319301fb04e2f40430206e913092ed54e2e030f845f843f842c8f841cf16cbffccccc9ed541501f682084c4b4001a013bef2e20801d3fffa4021d70b01c0009231029133e202d4d430f844f82870f842c8cbffc9c85003cf16cb07ccc97020c8cb0113f400f400cb00c920f9007074c8cb02ca07cbffc9d0f843d070c804d014cf16f843f842c8cbfff828cf16c903d430c8cc13ccc9c8cc17cbff5007cf1614cc15cc14013ccc43308040db3cf842a4f862f845f843f842c8f841cf16cbffccccc9ed54150078708010c8cb055006cf165004fa0214cb68216e947032cb009bc858cf17c97158cb00f400e2226e95327058cb0099c85003cf17c958f400e2c901fb003366cbbe'
codebase_version = 5
def __init__(self, **kwargs):
+28 -5
View File
@@ -12,11 +12,34 @@ kwargs = {}
if int(os.getenv('INIT_DEPLOY_PLATFORM_CONTRACT', 0)) == 0:
kwargs['address'] = Address(MY_PLATFORM_CONTRACT)
platform = Platform(
admin_address=Address('UQAjz4Kdqoo4_Obg-UrUmuhoUB2W00vngZoX0MnAAnetZuAk'),
def platform_with_salt(s: int = 0):
return Platform(
admin_address=Address('UQD3XALhbETNo7ItrdPNFzMJtRHC5u6dIb39DCYa40jnWZdg'),
blank_code=Cell.one_from_boc(Blank.code),
cop_code=Cell.one_from_boc(COP_NFT.code),
collection_content_uri=f'{PROJECT_HOST}/api/platform-metadata.json',
collection_content_uri=f'{PROJECT_HOST}/api/platform-metadata.json' + f"?s={s}",
**kwargs
)
)
platform = platform_with_salt()
if int(os.getenv('INIT_DEPLOY_PLATFORM_CONTRACT', 0)) == 1:
def is_nice_address(address: Address):
bounceable_addr = address.to_string(True, True, True)
non_bounceable_addr = address.to_string(True, True, False)
if '-' in bounceable_addr or '-' in non_bounceable_addr:
return False
if '_' in bounceable_addr or '_' in non_bounceable_addr:
return False
if bounceable_addr[-1] != 'A':
return False
return True
salt_value = 0
while not is_nice_address(platform.address):
platform = platform_with_salt(salt_value)
salt_value += 1
+9 -5
View File
@@ -7,7 +7,12 @@ load_dotenv(dotenv_path='.env')
PROJECT_HOST = os.getenv('PROJECT_HOST', 'http://127.0.0.1:8080')
SANIC_PORT = int(os.getenv('SANIC_PORT', '8080'))
# Path inside the running backend container where content files are visible
UPLOADS_DIR = os.getenv('UPLOADS_DIR', '/app/data')
# Host path where the same content directory is mounted (used for docker -v from within container)
BACKEND_DATA_DIR_HOST = os.getenv('BACKEND_DATA_DIR_HOST', '/Storage/storedContent')
# Host path for converter logs (used for docker -v). Optional.
BACKEND_LOGS_DIR_HOST = os.getenv('BACKEND_LOGS_DIR_HOST', '/Storage/logs/converter')
if not os.path.exists(UPLOADS_DIR):
os.makedirs(UPLOADS_DIR)
@@ -19,9 +24,8 @@ import httpx
TELEGRAM_BOT_USERNAME = httpx.get(f"https://api.telegram.org/bot{TELEGRAM_API_KEY}/getMe").json()['result']['username']
CLIENT_TELEGRAM_BOT_USERNAME = httpx.get(f"https://api.telegram.org/bot{CLIENT_TELEGRAM_API_KEY}/getMe").json()['result']['username']
MYSQL_URI = os.environ['MYSQL_URI']
MYSQL_DATABASE = os.environ['MYSQL_DATABASE']
# Unified database URL (PostgreSQL)
DATABASE_URL = os.environ['DATABASE_URL']
LOG_LEVEL = os.getenv('LOG_LEVEL', 'DEBUG')
LOG_DIR = os.getenv('LOG_DIR', 'logs')
@@ -32,7 +36,7 @@ _now_str = datetime.now().strftime("%Y-%m-%d_%H-%M-%S")
LOG_FILEPATH = f"{LOG_DIR}/{_now_str}.log"
WEB_APP_URLS = {
'uploadContent': f"https://web2-client.vercel.app/uploadContent"
'uploadContent': f"https://my-public-node-8.projscale.dev/uploadContent"
}
ALLOWED_CONTENT_TYPES = [
@@ -48,5 +52,5 @@ TONCENTER_HOST = os.getenv('TONCENTER_HOST', 'https://toncenter.com/api/v2/')
TONCENTER_API_KEY = os.getenv('TONCENTER_API_KEY')
TONCENTER_V3_HOST = os.getenv('TONCENTER_V3_HOST', 'https://toncenter.com/api/v3/')
MY_PLATFORM_CONTRACT = 'EQDmWp6hbJlYUrXZKb9N88sOrTit630ZuRijfYdXEHLtheMY'
MY_PLATFORM_CONTRACT = 'EQBVjuNuaIK87v9nm7mghgJ41ikqfx3GNBFz05GfmNbRQ9EA'
MY_FUND_ADDRESS = 'UQDarChHFMOI2On9IdHJNeEKttqepgo0AY4bG1trw8OAAwMY'
+27 -21
View File
@@ -36,9 +36,10 @@ async def create_new_encryption_key(db_session, user_id: int = None) -> KnownKey
meta={"I_user_id": user_id} if user_id else None,
created=datetime.now()
)
from sqlalchemy import select
db_session.add(new_key)
db_session.commit()
new_key = db_session.query(KnownKey).filter(KnownKey.seed_hash == new_seed_hash).first()
await db_session.commit()
new_key = (await db_session.execute(select(KnownKey).where(KnownKey.seed_hash == new_seed_hash))).scalars().first()
assert new_key, "Key not created"
return new_key
@@ -46,42 +47,51 @@ async def create_new_encryption_key(db_session, user_id: int = None) -> KnownKey
async def create_encrypted_content(
db_session, decrypted_content: StoredContent,
) -> StoredContent:
encrypted_content = db_session.query(StoredContent).filter(
StoredContent.id == decrypted_content.decrypted_content_id
).first()
from sqlalchemy import select
# Try to find an already created encrypted counterpart for this decrypted content
encrypted_content = (
await db_session.execute(
select(StoredContent).where(StoredContent.decrypted_content_id == decrypted_content.id)
)
).scalars().first()
if encrypted_content:
make_log("create_encrypted_content", f"(d={decrypted_content.cid.serialize_v2()}) => (e={encrypted_content.cid.serialize_v2()}): already exist (found by decrypted content)", level="debug")
return encrypted_content
encrypted_content = None
if decrypted_content.key is None:
# Avoid accessing relationship attributes in async context to prevent MissingGreenlet
if not decrypted_content.key_id:
key = await create_new_encryption_key(db_session, user_id=decrypted_content.user_id)
decrypted_content.key_id = key.id
db_session.commit()
decrypted_content = db_session.query(StoredContent).filter(
StoredContent.id == decrypted_content.id
).first()
await db_session.commit()
assert decrypted_content.key_id, "Key not assigned"
# Explicitly load the key to avoid lazy-loading via relationship in async mode
key = (
await db_session.execute(select(KnownKey).where(KnownKey.id == decrypted_content.key_id))
).scalars().first()
# If the referenced key is missing or malformed, create a fresh one
if not key or not key.seed:
key = await create_new_encryption_key(db_session, user_id=decrypted_content.user_id)
decrypted_content.key_id = key.id
await db_session.commit()
decrypted_path = os.path.join(UPLOADS_DIR, decrypted_content.hash)
decrypted_bin = b58decode(decrypted_content.hash)
key = decrypted_content.key
cipher = AESCipher(key.seed_bin)
encrypted_bin = cipher.encrypt(decrypted_bin)
encrypted_hash_bin = sha256(encrypted_bin).digest()
encrypted_hash = b58encode(encrypted_hash_bin).decode()
encrypted_content = db_session.query(StoredContent).filter(
StoredContent.hash == encrypted_hash
).first()
encrypted_content = (await db_session.execute(select(StoredContent).where(StoredContent.hash == encrypted_hash))).scalars().first()
if encrypted_content:
make_log("create_encrypted_content", f"(d={decrypted_content.cid.serialize_v2()}) => (e={encrypted_content.cid.serialize_v2()}): already exist (found by encrypted_hash)", level="debug")
return encrypted_content
encrypted_content = None
encrypted_meta = decrypted_content.meta
encrypted_meta = dict(decrypted_content.meta or {})
encrypted_meta["encrypt_algo"] = "AES256"
encrypted_content = StoredContent(
@@ -99,19 +109,15 @@ async def create_encrypted_content(
created=datetime.now(),
)
db_session.add(encrypted_content)
db_session.commit()
await db_session.commit()
encrypted_path = os.path.join(UPLOADS_DIR, encrypted_hash)
async with aiofiles.open(encrypted_path, mode='wb') as file:
await file.write(encrypted_bin)
encrypted_content = db_session.query(StoredContent).filter(
StoredContent.hash == encrypted_hash
).first()
encrypted_content = (await db_session.execute(select(StoredContent).where(StoredContent.hash == encrypted_hash))).scalars().first()
assert encrypted_content, "Content not created"
make_log("create_encrypted_content", f"(d={decrypted_content.cid.serialize_v2()}) => (e={encrypted_content.cid.serialize_v2()}): created new content/bin", level="debug")
return encrypted_content
+98 -27
View File
@@ -1,44 +1,115 @@
from os import getenv, urandom
import os
import time
import json
from nacl.bindings import crypto_sign_seed_keypair
from tonsdk.utils import Address
from app.core._blockchain.ton.wallet_v3cr3 import WalletV3CR3
from app.core.models._config import ServiceConfig
from app.core.storage import db_session
from app.core.logger import make_log
import os
from sqlalchemy import create_engine, inspect
from sqlalchemy.orm import Session
from typing import Optional
from app.core.models._config import ServiceConfigValue
def load_hot_pair():
with db_session() as session:
service_config = ServiceConfig(session)
hot_seed = service_config.get('private_key')
if hot_seed is None:
make_log("HotWallet", "No seed found, generating new one", level='info')
hot_seed = os.getenv("TON_INIT_HOT_SEED")
if not hot_seed:
hot_seed = urandom(32)
make_log("HotWallet", f"Generated random seed")
def _load_seed_from_env_or_generate() -> bytes:
seed_hex = os.getenv("TON_INIT_HOT_SEED")
if seed_hex:
make_log("HotWallet", "Loaded seed from env")
return bytes.fromhex(seed_hex)
make_log("HotWallet", "No seed provided; generating ephemeral seed", level='info')
return urandom(32)
def _init_seed_via_db() -> bytes:
"""Store and read hot seed from PostgreSQL service_config (key='private_key').
Primary node writes it once; workers wait until it appears.
"""
from app.core._config import DATABASE_URL
engine = create_engine(DATABASE_URL, pool_pre_ping=True)
role = os.getenv("NODE_ROLE", "worker").lower()
def db_ready(conn) -> bool:
try:
inspector = inspect(conn)
return inspector.has_table('service_config')
except Exception:
return False
# Wait for table to exist
start = time.time()
# Wait for table existence, reconnecting to avoid stale transactions
while True:
with engine.connect() as conn:
if db_ready(conn):
break
time.sleep(0.5)
if time.time() - start > 120:
raise TimeoutError("service_config table not available")
def read_seed() -> Optional[bytes]:
# Use a fresh connection/session per read to avoid snapshot staleness
try:
with engine.connect() as rconn:
with Session(bind=rconn) as s:
row = s.query(ServiceConfigValue).filter(ServiceConfigValue.key == 'private_key').first()
if not row:
return None
packed = row.packed_value or {}
if isinstance(packed, str):
packed = json.loads(packed)
seed_hex = packed.get('value')
return bytes.fromhex(seed_hex) if seed_hex else None
except Exception:
return None
seed = read_seed()
if seed:
return seed
if role == "primary":
seed = _load_seed_from_env_or_generate()
# Try insert; if another primary raced, ignore
try:
with engine.connect() as wconn:
with Session(bind=wconn) as s:
s.add(ServiceConfigValue(key='private_key', packed_value={"value": seed.hex()}))
s.commit()
make_log("HotWallet", "Seed saved in service_config by primary", level='info')
return seed
except Exception:
# Read again in case of race
seed2 = read_seed()
if seed2:
return seed2
raise
else:
hot_seed = bytes.fromhex(hot_seed)
make_log("HotWallet", f"Loaded seed from env")
service_config.set('private_key', hot_seed.hex())
return load_hot_pair()
hot_seed = bytes.fromhex(hot_seed)
public_key, private_key = crypto_sign_seed_keypair(hot_seed)
return hot_seed, public_key, private_key
make_log("HotWallet", "Worker waiting for seed in service_config...", level='info')
while True:
seed = read_seed()
if seed:
return seed
time.sleep(0.5)
_extra_ton_wallet_options = {}
if getenv('TON_CUSTOM_WALLET_ADDRESS'):
_extra_ton_wallet_options['address'] = Address(getenv('TON_CUSTOM_WALLET_ADDRESS'))
hot_seed, hot_pubkey, hot_privkey = load_hot_pair()
service_wallet = WalletV3CR3(
private_key=hot_privkey,
public_key=hot_pubkey,
def _init_wallet():
# Primary writes to DB; workers wait and read from DB
hot_seed_bytes = _init_seed_via_db()
pub, priv = crypto_sign_seed_keypair(hot_seed_bytes)
wallet = WalletV3CR3(
private_key=priv,
public_key=pub,
**_extra_ton_wallet_options
)
)
return hot_seed_bytes, pub, priv, wallet
hot_seed, hot_pubkey, hot_privkey, service_wallet = _init_wallet()
+8 -6
View File
@@ -1,10 +1,12 @@
from app.core.models import Asset
from sqlalchemy.ext.asyncio import AsyncEngine
from app.core.models import BlockchainTask
from app.core.models.base import AlchemyBase
def create_maria_tables(engine):
"""Create all tables in the database."""
Asset()
AlchemyBase.metadata.create_all(engine)
async def create_db_tables(engine: AsyncEngine):
"""Create all tables in the database (PostgreSQL, async)."""
# ensure model import side-effects initialize mappers
BlockchainTask()
async with engine.begin() as conn:
await conn.run_sync(AlchemyBase.metadata.create_all)
+2 -1
View File
@@ -5,7 +5,6 @@ from httpx import AsyncClient
from app.core._config import PROJECT_HOST
from app.core._crypto.signer import Signer
from app.core._secrets import hot_seed
from app.core.logger import make_log
@@ -17,6 +16,8 @@ async def send_status(service: str, status: str):
'status': status,
}
message_bytes = dumps(message).encode()
# Lazy import to avoid triggering _secrets before DB is ready
from app.core._secrets import hot_seed
signer = Signer(hot_seed)
message_signature = signer.sign(message_bytes)
async with AsyncClient() as client:
+3 -2
View File
@@ -56,9 +56,10 @@ class AuthenticationMixin:
},
created=datetime.fromtimestamp(init_ts)
)
from sqlalchemy import select
db_session.add(new_key)
db_session.commit()
new_key = db_session.query(KnownKey).filter(KnownKey.seed_hash == new_key.seed_hash).first()
await db_session.commit()
new_key = (await db_session.execute(select(KnownKey).where(KnownKey.seed_hash == new_key.seed_hash))).scalars().first()
assert new_key, "Key not created"
make_log("auth", f"[new-K] User {user_id} created new {token_type} key {new_key.id}")
return {
+129 -79
View File
@@ -4,8 +4,9 @@ import os
import uuid
import json
import shutil
import magic # python-magic for MIME detection
from base58 import b58decode, b58encode
from sqlalchemy import and_, or_
from sqlalchemy import and_, or_, select
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
@@ -13,14 +14,14 @@ 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._config import UPLOADS_DIR, BACKEND_DATA_DIR_HOST, BACKEND_LOGS_DIR_HOST
from app.core.content.content_id import ContentId
async def convert_loop(memory):
with db_session() as session:
async with db_session() as session:
# Query for unprocessed encrypted content
unprocessed_encrypted_content = session.query(StoredContent).filter(
unprocessed_encrypted_content = (await session.execute(select(StoredContent).where(
and_(
StoredContent.type == "onchain/content",
or_(
@@ -28,66 +29,116 @@ async def convert_loop(memory):
StoredContent.ipfs_cid == None,
)
)
).first()
))).scalars().first()
if not unprocessed_encrypted_content:
make_log("ConvertProcess", "No content to convert", level="debug")
return
# Достаем расшифрованный файл
decrypted_content = (await session.execute(select(StoredContent).where(
StoredContent.id == unprocessed_encrypted_content.decrypted_content_id
))).scalars().first()
if not decrypted_content:
make_log("ConvertProcess", "Decrypted content not found", level="error")
return
# Определяем путь и расширение входного файла
# Путь внутри текущего контейнера (доступен Python процессу)
input_file_container = os.path.join(UPLOADS_DIR, decrypted_content.hash)
# Хостовый путь (нужен для docker -v маппинга при запуске конвертера)
input_file_host = os.path.join(BACKEND_DATA_DIR_HOST, 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_container, 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']
}
}
await 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} with preview interval {preview_interval}", level="info")
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
make_log(
"ConvertProcess",
f"Processing content {unprocessed_encrypted_content.id} as {content_kind} with preview interval {preview_interval}",
level="info"
)
# List of conversion options to process
# Выбираем опции конвертации для видео и аудио
if content_kind == "video":
REQUIRED_CONVERT_OPTIONS = ['high', 'low', 'low_preview']
converted_content = {} # Mapping: option -> sha256 hash of output file
else:
REQUIRED_CONVERT_OPTIONS = ['high', 'low'] # no preview for audio
# Define input file path and extract its extension from filename
input_file_path = f"/Storage/storedContent/{decrypted_content.hash}"
converted_content = {}
# Директория логов на хосте для docker-контейнера конвертера
logs_dir_host = BACKEND_LOGS_DIR_HOST
input_ext = unprocessed_encrypted_content.filename.split('.')[-1] if '.' in unprocessed_encrypted_content.filename else "mp4"
# Logs directory mapping
logs_dir = "/Storage/logs/converter"
# Process each conversion option in sequence
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 # 'high' or 'low'
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}"
# Директория вывода в текущем контейнере (та же что и в UPLOADS_DIR, смонтирована с хоста)
output_dir_container = os.path.join(UPLOADS_DIR, "converter-output", output_uuid)
os.makedirs(output_dir_container, exist_ok=True)
# Соответствующая директория на хосте — нужна для docker -v
output_dir_host = os.path.join(BACKEND_DATA_DIR_HOST, "converter-output", output_uuid)
# Build the docker command with appropriate volume mounts and parameters
# 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",
# Важно: источники - это ХОСТОВЫЕ пути, так как docker демону они нужны на хосте
"-v", f"{input_file_host}:/app/input:ro",
"-v", f"{output_dir_host}:/app/output",
"-v", f"{logs_dir_host}:/app/logs",
"media_converter",
"--ext", input_ext,
"--quality", quality
]
if 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(
*cmd,
stdout=asyncio.subprocess.PIPE,
@@ -98,22 +149,21 @@ async def convert_loop(memory):
make_log("ConvertProcess", f"Docker conversion failed for option {option}: {stderr.decode()}", level="error")
return
# List files in the output directory
# List files in output dir
try:
files = os.listdir(output_dir.replace("/Storage/storedContent", "/app/data"))
files = os.listdir(output_dir_container)
except Exception as e:
make_log("ConvertProcess", f"Error reading output directory {output_dir}: {e}", level="error")
return
# Exclude 'output.json' and expect exactly one media output file
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])
output_file = os.path.join(output_dir_container, 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(
"sha256sum", output_file,
stdout=asyncio.subprocess.PIPE,
@@ -126,21 +176,18 @@ async def convert_loop(memory):
file_hash = hash_stdout.decode().split()[0]
file_hash = b58encode(bytes.fromhex(file_hash)).decode()
if not session.query(StoredContent).filter(
StoredContent.hash == file_hash
).first():
# Save new StoredContent if not exists
if not (await session.execute(select(StoredContent).where(StoredContent.hash == file_hash))).scalars().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,
},
meta={'encrypted_file_hash': unprocessed_encrypted_content.hash},
created=datetime.now(),
)
session.add(new_content)
session.commit()
await session.commit()
save_path = os.path.join(UPLOADS_DIR, file_hash)
try:
@@ -156,41 +203,29 @@ async def convert_loop(memory):
converted_content[option] = file_hash
# Process output.json: read its contents and update meta['ffprobe_meta']
output_json_path = os.path.join(output_dir.replace("/Storage/storedContent", "/app/data"), "output.json")
if os.path.exists(output_json_path):
if unprocessed_encrypted_content.meta.get('ffprobe_meta') is None:
# Process output.json for ffprobe_meta
output_json_path = os.path.join(output_dir_container, "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:
output_json_content = f.read()
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
ffprobe_meta = json.load(f)
unprocessed_encrypted_content.meta = {
**unprocessed_encrypted_content.meta,
'ffprobe_meta': ffprobe_meta
}
else:
make_log("ConvertProcess", f"output.json not found for option {option}", level="error")
# Remove the output directory after processing
try:
shutil.rmtree(output_dir.replace("/Storage/storedContent", "/app/data"))
except Exception as e:
make_log("ConvertProcess", f"Error removing output directory {output_dir}: {e}", level="error")
# Continue even if deletion fails
make_log("ConvertProcess", f"Error handling output.json for option {option}: {e}", level="error")
# Cleanup output directory
try:
shutil.rmtree(output_dir_container)
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'])
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'])
@@ -199,32 +234,47 @@ async def convert_loop(memory):
**unprocessed_encrypted_content.meta,
'converted_content': converted_content
}
await session.commit()
session.commit()
# Notify user if needed
if not unprocessed_encrypted_content.meta.get('upload_notify_msg_id'):
wallet_owner_connection = session.query(WalletConnection).filter(
wallet_owner_connection = (await session.execute(select(WalletConnection).where(
WalletConnection.wallet_address == unprocessed_encrypted_content.owner_address
).order_by(WalletConnection.id.desc()).first()
).order_by(WalletConnection.id.desc()))).scalars().first()
if wallet_owner_connection:
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)
unprocessed_encrypted_content.meta = {
**unprocessed_encrypted_content.meta,
'upload_notify_msg_id': await wallet_owner_bot.send_content(session, unprocessed_encrypted_content)
}
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)
await session.commit()
session.commit()
async def main_fn(memory):
make_log("ConvertProcess", "Service started", level="info")
seqno = 0
while True:
try:
make_log("ConvertProcess", "Service running", level="debug")
rid = __import__('uuid').uuid4().hex[:8]
try:
from app.core.log_context import ctx_rid
ctx_rid.set(rid)
except BaseException:
pass
make_log("ConvertProcess", "Service running", level="debug", rid=rid)
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")
make_log("ConvertProcess", f"Error: {e}", level="error", rid=locals().get('rid'))
await asyncio.sleep(3)
finally:
try:
from app.core.log_context import ctx_rid
ctx_rid.set(None)
except BaseException:
pass
+51 -35
View File
@@ -12,11 +12,13 @@ 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.user import User
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 sqlalchemy import select, and_, desc
from app.core.storage import db_session
import os
import traceback
@@ -33,7 +35,7 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
platform_found = True
make_log("Indexer", "Service running", level="debug")
with db_session() as session:
async with db_session() as session:
try:
result = await toncenter.run_get_method('EQD8TJ8xEWB1SpnRE4d89YO3jl0W0EiBnNS4IBaHaUmdfizE', 'get_pool_data')
assert result['exit_code'] == 0, f"Error in get-method: {result}"
@@ -41,42 +43,46 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
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()])
await 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(
from sqlalchemy import cast
from sqlalchemy.dialects.postgresql import JSONB
new_licenses = (await session.execute(select(UserContent).where(
and_(
~UserContent.meta.contains({'notification_sent': True}),
~(cast(UserContent.meta, JSONB).contains({'notification_sent': True})),
UserContent.type == 'nft/listen'
)
).all()
))).scalars().all()
for new_license in new_licenses:
licensed_content = session.query(StoredContent).filter(
licensed_content = (await session.execute(select(StoredContent).where(
StoredContent.id == new_license.content_id
).first()
))).scalars().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)
content_metadata = await licensed_content.metadata_json_async(session)
assert content_metadata, "No content metadata found"
if not (licensed_content.owner_address == new_license.owner_address):
try:
user = new_license.user
user = await session.get(User, new_license.user_id)
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
wallet_owner_connection = (await session.execute(
select(WalletConnection).where(
WalletConnection.wallet_address == licensed_content.owner_address,
WalletConnection.invalidated == False
).order_by(desc(WalletConnection.id))
)).scalars().first()
wallet_owner_user = await session.get(User, wallet_owner_connection.user_id) if wallet_owner_connection else None
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(
user.translated('p_licenseWasBought').format(
username=user.front_format(),
@@ -89,21 +95,19 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
make_log("IndexerSendNewLicense", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
new_license.meta = {**new_license.meta, 'notification_sent': True}
session.commit()
await session.commit()
content_without_cid = session.query(StoredContent).filter(
StoredContent.content_id == None
)
content_without_cid = (await session.execute(select(StoredContent).where(StoredContent.content_id == None))).scalars().all()
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()
await session.commit()
last_known_index_ = session.query(StoredContent).filter(
StoredContent.onchain_index != None
).order_by(StoredContent.onchain_index.desc()).first()
last_known_index_ = (await session.execute(
select(StoredContent).where(StoredContent.onchain_index != None).order_by(StoredContent.onchain_index.desc())
)).scalars().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")
@@ -196,14 +200,13 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
user_wallet_connection = None
if item_owner_address:
user_wallet_connection = session.query(WalletConnection).filter(
user_wallet_connection = (await session.execute(select(WalletConnection).where(
WalletConnection.wallet_address == item_owner_address.to_string(1, 1, 1)
).first()
))).scalars().first()
encrypted_stored_content = session.query(StoredContent).filter(
StoredContent.hash == item_content_hash_str,
# StoredContent.type.like("local%")
).first()
encrypted_stored_content = (await session.execute(select(StoredContent).where(
StoredContent.hash == item_content_hash_str
))).scalars().first()
if encrypted_stored_content:
is_duplicate = encrypted_stored_content.type.startswith("onchain") \
and encrypted_stored_content.onchain_index != item_index
@@ -215,7 +218,7 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
user = None
if user_wallet_connection:
encrypted_stored_content.user_id = user_wallet_connection.user_id
user = user_wallet_connection.user
user = await session.get(User, user_wallet_connection.user_id)
if user:
user_uploader_wrapper = Wrapped_CBotChat(memory._telegram_bot, chat_id=user.telegram_id, user=user, db_session=session)
@@ -234,14 +237,15 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
)
try:
for hint_message in session.query(KnownTelegramMessage).filter(
result = await session.execute(select(KnownTelegramMessage).where(
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():
))
for hint_message in result.scalars().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")
@@ -260,7 +264,7 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
**item_metadata_packed
}
session.commit()
await session.commit()
return platform_found, seqno
else:
item_metadata_packed['copied_from'] = encrypted_stored_content.id
@@ -282,7 +286,7 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
updated=datetime.now()
)
session.add(onchain_stored_content)
session.commit()
await session.commit()
make_log("Indexer", f"Item indexed: {item_content_hash_str}", level="info")
last_known_index += 1
@@ -295,15 +299,27 @@ async def main_fn(memory, ):
seqno = 0
while True:
try:
rid = __import__('uuid').uuid4().hex[:8]
try:
from app.core.log_context import ctx_rid
ctx_rid.set(rid)
except BaseException:
pass
make_log("Indexer", f"Loop start", level="debug", rid=rid)
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")
make_log("Indexer", f"Error: {e}" + '\n' + traceback.format_exc(), level="error", rid=locals().get('rid'))
if platform_found:
await send_status("indexer", f"working (seqno={seqno})")
await asyncio.sleep(5)
seqno += 1
try:
from app.core.log_context import ctx_rid
ctx_rid.set(None)
except BaseException:
pass
+33 -18
View File
@@ -3,7 +3,7 @@ from base64 import b64decode
from datetime import datetime, timedelta
from base58 import b58encode
from sqlalchemy import and_
from sqlalchemy import and_, or_, select, desc
from tonsdk.boc import Cell
from tonsdk.utils import Address
@@ -27,7 +27,7 @@ import traceback
async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
make_log("LicenseIndex", "Service running", level="debug")
with db_session() as session:
async with db_session() as session:
async def check_telegram_stars_transactions():
# Проверка звездных telegram транзакций, обновление paid
offset = {'desc': 'Статичное число заранее известного количества транзакций, которое даже не знает наш бот', 'value': 1}['value'] + \
@@ -45,19 +45,19 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
continue
try:
existing_invoice = session.query(StarsInvoice).filter(
existing_invoice = (await session.execute(select(StarsInvoice).where(
StarsInvoice.external_id == star_payment.source.invoice_payload
).first()
))).scalars().first()
if not existing_invoice:
continue
if star_payment.amount == existing_invoice.amount:
if not existing_invoice.paid:
existing_invoice.paid = True
session.commit()
await session.commit()
licensed_content = session.query(StoredContent).filter(StoredContent.hash == existing_invoice.content_hash).first()
user = session.query(User).filter(User.id == existing_invoice.user_id).first()
licensed_content = (await session.execute(select(StoredContent).where(StoredContent.hash == existing_invoice.content_hash))).scalars().first()
user = (await session.execute(select(User).where(User.id == existing_invoice.user_id))).scalars().first()
await (Wrapped_CBotChat(memory._client_telegram_bot, chat_id=user.telegram_id, user=user, db_session=session)).send_content(
session, licensed_content
@@ -73,33 +73,36 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
make_log("StarsProcessing", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
# Проверка кошельков пользователей на появление новых NFT, добавление их в базу как неопознанные
for user in session.query(User).filter(
User.last_use > datetime.now() - timedelta(minutes=10)
).all():
if not user.wallet_address(session):
make_log("LicenseIndex", f"User {user.id} has no wallet address", level="info")
users = (await session.execute(select(User).where(
User.last_use > datetime.now() - timedelta(hours=4)
).order_by(User.updated.asc()))).scalars().all()
for user in users:
user_wallet_address = await user.wallet_address_async(session)
if not user_wallet_address:
make_log("LicenseIndex", f"User {user.id} has no wallet address", level="debug")
continue
make_log("LicenseIndex", f"User {user.id} has wallet address {user_wallet_address}", level="debug")
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)
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="debug")
if must_skip:
continue
try:
await user.scan_owned_user_content(session)
user.meta = {**user.meta, 'last_updated_licenses': datetime.now().isoformat()}
session.commit()
await session.commit()
except BaseException as e:
make_log("LicenseIndex", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
# Проверка NFT на актуальность данных, в том числе уже проверенные
process_content = session.query(UserContent).filter(
process_content = (await session.execute(select(UserContent).where(
and_(
UserContent.type.startswith('nft/'),
UserContent.updated < (datetime.now() - timedelta(minutes=60)),
)
).first()
).order_by(UserContent.updated.asc()))).scalars().first()
if process_content:
make_log("LicenseIndex", f"Syncing content with blockchain: {process_content.id}", level="info")
try:
@@ -108,7 +111,7 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
make_log("LicenseIndex", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
finally:
process_content.updated = datetime.now()
session.commit()
await session.commit()
return platform_found, seqno
@@ -119,14 +122,26 @@ async def main_fn(memory, ):
seqno = 0
while True:
try:
rid = __import__('uuid').uuid4().hex[:8]
try:
from app.core.log_context import ctx_rid
ctx_rid.set(rid)
except BaseException:
pass
make_log("LicenseIndex", f"Loop start", level="debug", rid=rid)
platform_found, seqno = await license_index_loop(memory, platform_found, seqno)
if platform_found:
await send_status("licenses", f"working (seqno={seqno})")
except BaseException as e:
make_log("LicenseIndex", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
make_log("LicenseIndex", f"Error: {e}" + '\n' + traceback.format_exc(), level="error", rid=locals().get('rid'))
await asyncio.sleep(1)
seqno += 1
try:
from app.core.log_context import ctx_rid
ctx_rid.set(None)
except BaseException:
pass
# if __name__ == '__main__':
# loop = asyncio.get_event_loop()
+211 -9
View File
@@ -1,13 +1,22 @@
import asyncio
from base64 import b64decode
import os
import traceback
import httpx
from tonsdk.boc import begin_cell
from sqlalchemy import and_, func
from tonsdk.boc import begin_cell, Cell
from tonsdk.contract.wallet import Wallets
from tonsdk.utils import HighloadQueryId
from datetime import datetime, timedelta
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
@@ -75,26 +84,219 @@ async def main_fn(memory):
await asyncio.sleep(15)
return await main_fn(memory)
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.02 * 10 ** 9),
'send_mode': 1,
'payload': begin_cell().store_uint(0, 32).end_cell()
}], sw_seqno_value
)['message'].to_boc(False)
)
await send_status("ton_daemon", "working: topup highload wallet")
await asyncio.sleep(15)
return await main_fn(memory)
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:
try:
rid = __import__('uuid').uuid4().hex[:8]
try:
from app.core.log_context import ctx_rid
ctx_rid.set(rid)
except BaseException:
pass
sw_seqno_value = await get_sw_seqno()
make_log("TON", f"Service running ({sw_seqno_value})", level="debug")
make_log("TON", f"Service running ({sw_seqno_value})", level="debug", rid=rid)
# with db_session() as session:
# for stored_content in session.query(StoredContent).filter(StoredContent.uploaded == False).all():
# pass
async 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()
from sqlalchemy import select
in_msg_blockchain_task = (
await session.execute(
select(BlockchainTask).where(
and_(
BlockchainTask.seqno == in_msg_seqno,
BlockchainTask.epoch == in_msg_epoch,
)
)
)
).scalars().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
await 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', rid=rid)
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", rid=rid)
except BaseException as e:
make_log("TON_Daemon", f"Error while getting service wallet transactions: {e}", level="ERROR", rid=rid)
await send_status("ton_daemon", f"working: processing out-txs (seqno={sw_seqno_value})")
# Отправка подписанных сообщений
from sqlalchemy import select
_processing = (await session.execute(select(BlockchainTask).where(
BlockchainTask.status == 'processing'
).order_by(BlockchainTask.updated.asc()))).scalars().all()
for blockchain_task in _processing:
make_log("TON_Daemon", f"Processing task (processing) {blockchain_task.id}", rid=rid)
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'
await session.commit()
continue
await asyncio.sleep(0.5)
await send_status("ton_daemon", f"working: creating new messages (seqno={sw_seqno_value})")
# Создание новых подписей
_waiting = (await session.execute(select(BlockchainTask).where(BlockchainTask.status == 'wait'))).scalars().all()
for blockchain_task in _waiting:
try:
# Check processing tasks in current epoch < 3_000_000
from sqlalchemy import func
_cnt = (await session.execute(select(func.count()).select_from(BlockchainTask).where(
BlockchainTask.epoch == blockchain_task.epoch
))).scalar() or 0
if _cnt > 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))
from sqlalchemy import func
max_epoch_seqno = (
(await session.execute(select(func.max(BlockchainTask.seqno)).where(
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", rid=rid)
query_boc = begin_cell().end_cell().to_boc(False)
blockchain_task.meta = {
**blockchain_task.meta,
'sign_created': sign_created,
'signed_message': query_boc.hex(),
}
await session.commit()
make_log("TON", f"Created signed message for task {blockchain_task.id}" + '\n' + traceback.format_exc(), level="info", rid=rid)
except BaseException as e:
make_log("TON", f"Error processing task {blockchain_task.id}: {e}" + '\n' + traceback.format_exc(), level="error", rid=rid)
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")
make_log("TON", f"Error: {e}", level="error", rid=locals().get('rid'))
await asyncio.sleep(3)
finally:
try:
from app.core.log_context import ctx_rid
ctx_rid.set(None)
except BaseException:
pass
# if __name__ == '__main__':
# loop = asyncio.get_event_loop()
# loop.run_until_complete(main())
# loop.close()
+14 -4
View File
@@ -13,14 +13,26 @@ async def main_fn(memory):
seqno = 0
while True:
try:
make_log("Uploader", "Service running", level="debug")
rid = __import__('uuid').uuid4().hex[:8]
try:
from app.core.log_context import ctx_rid
ctx_rid.set(rid)
except BaseException:
pass
make_log("Uploader", f"Service running", level="debug", rid=rid)
await uploader_loop()
await asyncio.sleep(5)
await send_status("uploader_daemon", f"working (seqno={seqno})")
seqno += 1
except BaseException as e:
make_log("Uploader", f"Error: {e}", level="error")
make_log("Uploader", f"Error: {e}", level="error", rid=locals().get('rid'))
await asyncio.sleep(3)
finally:
try:
from app.core.log_context import ctx_rid
ctx_rid.set(None)
except BaseException:
pass
# if __name__ == '__main__':
# loop = asyncio.get_event_loop()
@@ -28,5 +40,3 @@ async def main_fn(memory):
# loop.close()
+108 -11
View File
@@ -1,8 +1,11 @@
import json
import asyncio
import os
import string
import aiofiles
from hashlib import sha256
import re # Added import
import unicodedata # Added import
from base58 import b58encode
from datetime import datetime, timedelta
@@ -23,7 +26,9 @@ async def create_new_content(
content_hash_bin = sha256(content_bin).digest()
content_hash_b58 = b58encode(content_hash_bin).decode()
new_content = db_session.query(StoredContent).filter(StoredContent.hash == content_hash_b58).first()
from sqlalchemy import select
result = await db_session.execute(select(StoredContent).where(StoredContent.hash == content_hash_b58))
new_content = result.scalars().first()
if new_content:
return new_content, False
@@ -35,8 +40,9 @@ async def create_new_content(
)
db_session.add(new_content)
db_session.commit()
new_content = db_session.query(StoredContent).filter(StoredContent.hash == content_hash_b58).first()
await db_session.commit()
result = await db_session.execute(select(StoredContent).where(StoredContent.hash == content_hash_b58))
new_content = result.scalars().first()
assert new_content, "Content not created (through utils)"
content_filepath = os.path.join(UPLOADS_DIR, content_hash_b58)
async with aiofiles.open(content_filepath, 'wb') as file:
@@ -45,34 +51,125 @@ async def create_new_content(
return new_content, True
# New helper functions for string cleaning
def _remove_emojis(text: str) -> str:
"""Removes common emoji characters from a string."""
# This regex covers many common emojis but might not be exhaustive.
emoji_pattern = re.compile(
"["
"\U0001F600-\U0001F64F" # emoticons
"\U0001F300-\U0001F5FF" # symbols & pictographs
"\U0001F680-\U0001F6FF" # transport & map symbols
"\U0001F1E0-\U0001F1FF" # flags (iOS)
"\U00002702-\U000027B0" # Dingbats
"\U000024C2-\U0001F251" # Various symbols
"]+",
flags=re.UNICODE,
)
return emoji_pattern.sub(r'', text)
def _clean_text_content(text: str, is_hashtag: bool = False) -> str:
"""
Cleans a string by removing emojis and unusual characters.
Level 1: Emoji removal.
Level 2: Unusual character cleaning (specific logic for hashtags).
"""
if not isinstance(text, str):
return ""
# Level 1: Remove emojis
text_no_emojis = _remove_emojis(text)
# Level 2: Clean unusual characters
if is_hashtag:
# Convert to lowercase
processed_text = text_no_emojis.lower()
# Replace hyphens, dots, spaces (and sequences) with a single underscore
processed_text = re.sub(r'[\s.-]+', '_', processed_text)
# Keep only lowercase letters (a-z), digits (0-9), and underscores
cleaned_text = re.sub(r'[^a-z0-9_]', '', processed_text)
# Remove leading/trailing underscores
cleaned_text = cleaned_text.strip('_')
# Consolidate multiple underscores into one
cleaned_text = re.sub(r'_+', '_', cleaned_text)
return cleaned_text
else: # For title, authors, or general text
# Normalize Unicode characters (e.g., NFKD form)
nfkd_form = unicodedata.normalize('NFKD', text_no_emojis)
# Keep letters (Unicode), numbers (Unicode), spaces, and basic punctuation
# This allows for a wider range of characters suitable for titles/names.
cleaned_text_chars = []
for char_in_nfkd in nfkd_form:
if not unicodedata.combining(char_in_nfkd): # remove combining diacritics
# Keep letters, numbers, spaces, and specific punctuation
cat = unicodedata.category(char_in_nfkd)
if cat.startswith('L') or cat.startswith('N') or cat.startswith('Z') or char_in_nfkd in '.,!?-':
cleaned_text_chars.append(char_in_nfkd)
cleaned_text = "".join(cleaned_text_chars)
# Normalize multiple spaces to a single space and strip leading/trailing spaces
cleaned_text = re.sub(r'\s+', ' ', cleaned_text).strip()
return cleaned_text
async def create_metadata_for_item(
db_session,
title: str = None,
cover_url: str = None,
authors: list = None,
hashtags: list = [],
downloadable: bool = False,
) -> StoredContent:
assert title, "No title provided"
# assert cover_url, "No cover_url provided"
assert len(title) > 3, "Title too short"
title = title[:100].strip()
# assert cover_url, "No cover_url provided" # Original comment, kept as is
# Clean title using the new helper function
cleaned_title = _clean_text_content(title, is_hashtag=False)
cleaned_title = cleaned_title[:100].strip() # Truncate and strip after cleaning
assert len(cleaned_title) > 3, f"Cleaned title '{cleaned_title}' (from original '{title}') is too short or became empty after cleaning."
# Process and clean hashtags
processed_hashtags = []
if hashtags and isinstance(hashtags, list):
for _h_tag_text in hashtags:
if isinstance(_h_tag_text, str):
cleaned_h = _clean_text_content(_h_tag_text, is_hashtag=True)
if cleaned_h: # Add only if not empty after cleaning
processed_hashtags.append(cleaned_h)
# Ensure uniqueness of hashtags and limit their count (e.g., to first 10 unique)
# Using dict.fromkeys to preserve order while ensuring uniqueness
processed_hashtags = list(dict.fromkeys(processed_hashtags))[:10]
item_metadata = {
'name': title,
'description': ' '.join([f"#{_h.replace(' ', '_')}" for _h in hashtags]),
'name': cleaned_title,
'attributes': [
# {
# 'trait_type': 'Artist',
# 'value': 'Unknown'
# },
],
'downloadable': downloadable,
'tags': processed_hashtags, # New field for storing the list of cleaned hashtags
}
# Generate description from the processed hashtags
item_metadata['description'] = ' '.join([f"#{h}" for h in processed_hashtags if h])
if cover_url:
item_metadata['image'] = cover_url
item_metadata['authors'] = [
''.join([_a_ch for _a_ch in _a if len(_a_ch.encode()) == 1]) for _a in (authors or [])[:500]
]
# Clean authors
cleaned_authors = []
if authors and isinstance(authors, list):
for author_name in (authors or [])[:500]: # Limit number of authors
if isinstance(author_name, str):
# Apply general cleaning to author names
# This replaces the old logic: ''.join([_a_ch for _a_ch in _a if len(_a_ch.encode()) == 1])
cleaned_author = _clean_text_content(author_name, is_hashtag=False)
if cleaned_author.strip(): # Ensure not empty
cleaned_authors.append(cleaned_author.strip()[:100]) # Limit length of each author name
item_metadata['authors'] = cleaned_authors
# Upload file
metadata_bin = json.dumps(item_metadata).encode()
+11
View File
@@ -0,0 +1,11 @@
from contextvars import ContextVar
# Correlation for HTTP requests
ctx_session_id = ContextVar('ctx_session_id', default=None)
ctx_user_id = ContextVar('ctx_user_id', default=None)
ctx_method = ContextVar('ctx_method', default=None)
ctx_path = ContextVar('ctx_path', default=None)
ctx_remote = ContextVar('ctx_remote', default=None)
# Correlation for background loop iterations
ctx_rid = ContextVar('ctx_rid', default=None)
+18
View File
@@ -1,5 +1,6 @@
import os
from app.core.projscale_logger import logger
from app.core.log_context import ctx_session_id, ctx_user_id, ctx_method, ctx_path, ctx_remote, ctx_rid
import logging
LOG_LEVELS = {
@@ -17,6 +18,23 @@ def make_log(issuer, message, *args, level='INFO', **kwargs):
assert level.upper() in LOG_LEVELS.keys(), f"Unknown log level"
_log = getattr(logger, level.lower())
# Merge context variables if not explicitly provided
context_fields = {
'sid': kwargs.get('sid') or kwargs.get('session_id') or ctx_session_id.get(),
'user_id': kwargs.get('user_id') or ctx_user_id.get(),
'method': kwargs.get('method') or ctx_method.get(),
'path': kwargs.get('path') or ctx_path.get(),
'remote': kwargs.get('remote') or ctx_remote.get(),
'rid': kwargs.get('rid') or ctx_rid.get(),
}
# Only include non-empty context
for k, v in list(context_fields.items()):
if v is None:
context_fields.pop(k)
# Do not override provided kwargs; merge missing only
for k, v in context_fields.items():
kwargs.setdefault(k, v)
log_buffer = f"[{issuer if not (issuer is None) else 'System'}] {message}"
if args:
log_buffer += f" | {args}"
+2
View File
@@ -11,3 +11,5 @@ from app.core.models.content.user_content import UserContent, UserAction
from app.core.models._config import ServiceConfigValue, ServiceConfig
from app.core.models.asset import Asset
from app.core.models.my_network import KnownNode, KnownNodeIncident, RemoteContentIndex
from app.core.models.promo import PromoAction
from app.core.models.tasks import BlockchainTask
+13 -15
View File
@@ -1,6 +1,6 @@
from app.core.models.base import AlchemyBase
from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime, JSON, Boolean
from sqlalchemy import Column, Integer, String, JSON, select
class ServiceConfigValue(AlchemyBase):
@@ -8,7 +8,7 @@ class ServiceConfigValue(AlchemyBase):
id = Column(Integer, autoincrement=True, primary_key=True)
key = Column(String(128), nullable=False, unique=True)
packed_value = Column(JSON, nullable=False, default={})
packed_value = Column(JSON, nullable=False, default=dict)
@property
def value(self):
@@ -19,20 +19,18 @@ class ServiceConfig:
def __init__(self, session):
self.session = session
def get(self, key, default=None):
result = self.session.query(ServiceConfigValue).filter(ServiceConfigValue.key == key).first()
async def get(self, key, default=None):
result = (await self.session.execute(select(ServiceConfigValue).where(ServiceConfigValue.key == key))).scalars().first()
return (result.value if result else None) or default
def set(self, key, value):
config_value = self.session.query(ServiceConfigValue).filter(
ServiceConfigValue.key == key
).first()
if not config_value:
config_value = ServiceConfigValue(key=key)
self.session.add(config_value)
self.session.commit()
return self.set(key, value)
async def set(self, key, value):
result = (await self.session.execute(select(ServiceConfigValue).where(ServiceConfigValue.key == key))).scalars().first()
if not result:
result = ServiceConfigValue(key=key)
self.session.add(result)
await self.session.commit()
return await self.set(key, value)
config_value.packed_value = {'value': value}
self.session.commit()
result.packed_value = {'value': value}
await self.session.commit()
return
+26 -144
View File
@@ -1,4 +1,4 @@
from sqlalchemy import and_
from sqlalchemy import and_, select
from app.core.models.node_storage import StoredContent
from app.core.models.content.user_content import UserContent, UserAction
from app.core.logger import make_log
@@ -27,17 +27,17 @@ class PlayerTemplates:
if not content.encrypted:
local_content = content
else:
local_content = db_session.query(StoredContent).filter_by(
id=content.decrypted_content_id
).first()
local_content = (await db_session.execute(select(StoredContent).where(StoredContent.id == content.decrypted_content_id))).scalars().first()
# TODO: add check decrypted_content by .format_json()['content_cid']
if local_content:
cd_log += f"Decrypted: {local_content.hash}. "
else:
cd_log += "Can't decrypt content. "
user_wallet_address = self.user.wallet_address(self.db_session)
user_existing_license = self.db_session.query(UserContent).filter_by(user_id=self.user.id, content_id=content.id).first()
user_wallet_address = await self.user.wallet_address_async(self.db_session)
user_existing_license = (await self.db_session.execute(select(UserContent).where(
and_(UserContent.user_id == self.user.id, UserContent.content_id == content.id)
))).scalars().first()
if local_content:
content_meta = content.json_format()
@@ -48,30 +48,17 @@ class PlayerTemplates:
except:
content_type, content_encoding = 'application', 'x-binary'
content_metadata = StoredContent.from_cid(db_session, content_meta.get('metadata_cid') or None)
content_metadata = await StoredContent.from_cid_async(db_session, content_meta.get('metadata_cid') or None)
with open(content_metadata.filepath, 'r') as f:
content_metadata_json = json.loads(f.read())
try:
cover_content = StoredContent.from_cid(self.db_session, content_meta.get('cover_cid') or None)
cover_content = await StoredContent.from_cid_async(self.db_session, content_meta.get('cover_cid') or None)
cd_log += f"Cover content: {cover_content.cid.serialize_v2()}. "
except BaseException as e:
cd_log += f"Can't get cover content: {e}. "
cover_content = None
local_content.meta['cover_cid'] = cover_content.cid.serialize_v2() if cover_content else None
local_content_cid = local_content.cid
local_content_url = f"{PROJECT_HOST}/api/v1.5/storage/{local_content_cid.serialize_v2()}"
converted_content = content.meta.get('converted_content')
if not converted_content:
r = await tg_process_template(
self, self.user.translated('p_playerContext_contentNotReady'),
message_id=message_id,
message_type='common'
)
return r
content_share_link = {
'text': self.user.translated('p_shareLinkContext').format(title=content_metadata_json.get('name', "")),
'url': f"https://t.me/{CLIENT_TELEGRAM_BOT_USERNAME}/content?startapp={content.cid.serialize_v2()}"
@@ -79,87 +66,8 @@ class PlayerTemplates:
if user_existing_license:
content_share_link['url'] = f"https://t.me/{CLIENT_TELEGRAM_BOT_USERNAME}/content?startapp={user_existing_license.onchain_address}"
preview_content = db_session.query(StoredContent).filter(
StoredContent.hash == converted_content['low_preview']
).first()
if preview_content.filename.split('.')[-1] in ['mov', 'mp4']:
content_type = 'video'
local_content_preview_url = preview_content.web_url
if content_type == 'audio':
audio_title = content_metadata_json.get('name', "").split(' - ')
if len(audio_title) > 1:
template_kwargs['performer'] = audio_title[0].strip()
audio_title = audio_title[1:]
template_kwargs['title'] = audio_title[0].strip()
template_kwargs['protect_content'] = True
template_kwargs['audio'] = URLInputFile(local_content_preview_url)
if cover_content:
template_kwargs['thumbnail'] = URLInputFile(cover_content.web_url)
if self.bot_id == 1:
inline_keyboard_array.append([
{
'text': self.user.translated('shareTrack_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:
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/{content_meta['item_address']}"
}])
elif content_type == 'video':
# Processing video
video_title = content_metadata_json.get('name', "")
template_kwargs['video'] = URLInputFile(local_content_preview_url)
template_kwargs['protect_content'] = True
if cover_content:
# Add thumbnail if cover content is available
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:
local_content = None
template_kwargs['photo'] = URLInputFile(cover_content.web_url)
if not local_content:
text = self.user.translated('p_playerContext_unsupportedContent').format(
@@ -169,37 +77,26 @@ class PlayerTemplates:
inline_keyboard_array = []
extra_buttons = []
else:
text = content_metadata_json.get('description').strip()
content_hashtags = content_metadata_json.get('description').strip()
if content_hashtags:
content_hashtags += '\n'
have_access = (
(content.owner_address == user_wallet_address)
or bool(user_existing_license)
or bool(self.db_session.query(StarsInvoice).filter(
and_(
StarsInvoice.user_id == self.user.id,
StarsInvoice.content_hash == content.hash,
StarsInvoice.paid == True
)
).first())
)
if have_access:
full_content = self.db_session.query(StoredContent).filter_by(
hash=content.meta.get('converted_content', {}).get('low') # TODO: support high quality
).first()
if content_type == 'audio':
# Restrict audio to 30 seconds if user does not have access
template_kwargs['audio'] = URLInputFile(full_content.web_url)
elif content_type == 'video':
# Restrict video to 30 seconds if user does not have access
template_kwargs['video'] = URLInputFile(full_content.web_url)
text = f"""<b>{content_metadata_json.get('name', 'Unnamed')}</b>
{content_hashtags}
Этот контент был загружен в MY
\t/ p2p content market /
<blockquote><a href="{content_share_link['url']}">🔴 «открыть в MY»</a></blockquote>"""
make_log("TG-Player", f"Send content {content_type} ({content_encoding}) to chat {self._chat_id}. {cd_log}")
for kmsg in self.db_session.query(KnownTelegramMessage).filter_by(
content_id=content.id,
chat_id=self._chat_id,
type=f'content/{content_type}',
deleted=False
).all():
kmsgs = (await self.db_session.execute(select(KnownTelegramMessage).where(
and_(
KnownTelegramMessage.content_id == content.id,
KnownTelegramMessage.chat_id == self._chat_id,
KnownTelegramMessage.type == f'content/{content_type}',
KnownTelegramMessage.deleted == False
)
))).scalars().all()
for kmsg in kmsgs:
await self.delete_message(kmsg.message_id)
r = await tg_process_template(
@@ -210,19 +107,4 @@ class PlayerTemplates:
content_id=content.id if content else None
)
if self.bot_id == 1:
if content.type == 'onchain/content':
if content_type == 'audio':
content.meta = {
**content.meta,
'telegram_file_cache': r.audio.file_id,
}
elif content_type == 'video':
content.meta = {
**content.meta,
'telegram_file_cache': r.video.file_id,
}
self.db_session.commit()
return r
+10 -7
View File
@@ -1,7 +1,7 @@
from aiogram import Bot, types
from datetime import datetime, timedelta
from sqlalchemy import and_
from sqlalchemy import and_, select
from app.core.logger import make_log
from app.core.models.messages import KnownTelegramMessage
@@ -46,14 +46,15 @@ class Wrapped_CBotChat(T, PlayerTemplates):
if self.db_session:
if message_type == 'common':
ci = 0
for oc_msg in self.db_session.query(KnownTelegramMessage).filter(
result = await self.db_session.execute(select(KnownTelegramMessage).where(
and_(
KnownTelegramMessage.type == 'common',
KnownTelegramMessage.bot_id == self.bot_id,
KnownTelegramMessage.chat_id == self._chat_id,
KnownTelegramMessage.deleted == False
)
).all():
))
for oc_msg in result.scalars().all():
make_log(self, f"Delete old message {oc_msg.message_id} {oc_msg.type} {oc_msg.bot_id} {oc_msg.chat_id}")
await self.delete_message(oc_msg.message_id)
ci += 1
@@ -75,7 +76,7 @@ class Wrapped_CBotChat(T, PlayerTemplates):
content_id=content_id
)
)
self.db_session.commit()
await self.db_session.commit()
else:
make_log(self, f"Unknown result type: {type(result)}", level='warning')
@@ -127,14 +128,16 @@ class Wrapped_CBotChat(T, PlayerTemplates):
message_id
)):
if self.db_session:
known_message = self.db_session.query(KnownTelegramMessage).filter(
known_message = (await self.db_session.execute(select(KnownTelegramMessage).where(
and_(
KnownTelegramMessage.bot_id == self.bot_id,
KnownTelegramMessage.chat_id == self._chat_id,
KnownTelegramMessage.message_id == message_id
).first()
)
))).scalars().first()
if known_message:
known_message.deleted = True
self.db_session.commit()
await self.db_session.commit()
except Exception as e:
make_log(self, f"Error deleting message {self._chat_id}/{message_id}. Error: {e}", level='warning')
return None
+14 -14
View File
@@ -1,5 +1,6 @@
from sqlalchemy import Column, Integer, String, DateTime, JSON, Boolean
from sqlalchemy.orm import relationship
from datetime import datetime
from app.core._defaults import DEFAULT_ASSET_INITOBJ
from app.core.models.base import AlchemyBase
@@ -15,10 +16,10 @@ class Asset(AlchemyBase):
network = Column(String(32), nullable=True)
address = Column(String(1024), nullable=True)
meta = Column(JSON, nullable=False, default={})
rates = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default=dict)
rates = Column(JSON, nullable=False, default=dict)
created = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
is_active = Column(Boolean, nullable=False, default=True)
balances = relationship('UserBalance', back_populates='asset')
@@ -29,22 +30,21 @@ class Asset(AlchemyBase):
AlchemyBase.metadata.create_all(engine)
@classmethod
def find(cls, session, **kwargs):
async def find_async(cls, session, **kwargs):
from sqlalchemy import select, func
if 'symbol' in kwargs:
kwargs['symbol'] = kwargs['symbol'].upper()
result = session.query(cls).filter_by(**kwargs)
results_count = result.count()
if results_count == 0:
any_count = session.query(cls).count()
result = await session.execute(select(cls).filter_by(**kwargs))
row = result.scalars().first()
if row:
return row
any_count = (await session.execute(select(func.count()).select_from(cls))).scalar() or 0
if any_count == 0:
init_asset = cls(**DEFAULT_ASSET_INITOBJ)
session.add(init_asset)
session.commit()
return cls.find(session, **kwargs)
await session.commit()
return await cls.find_async(session, **kwargs)
raise Exception(f"Asset not found: {kwargs}")
elif results_count == 1:
return result.first()
else:
raise Exception(f"Multiple assets found: {results_count}")
+6 -11
View File
@@ -1,7 +1,7 @@
import traceback
import base58
from sqlalchemy import and_
from sqlalchemy import and_, select
from app.core.logger import make_log
from app.core.models import StoredContent
@@ -57,13 +57,9 @@ class UserContentIndexationMixin:
values_slice = cc_indexator_data['values'].begin_parse()
content_hash_b58 = base58.b58encode(bytes.fromhex(hex(values_slice.read_uint(256))[2:])).decode()
make_log("UserContent", f"License ({self.onchain_address}) content hash: {content_hash_b58}", level="info")
stored_content = db_session.query(StoredContent).filter(
and_(
StoredContent.type == 'onchain/content',
StoredContent.hash == content_hash_b58,
)
).first()
stored_content = (await db_session.execute(select(StoredContent).where(
and_(StoredContent.type == 'onchain/content', StoredContent.hash == content_hash_b58)
))).scalars().first()
trusted_cop_address_result = await toncenter.run_get_method(stored_content.meta['item_address'], 'get_nft_address_by_index', [['num', cc_indexator_data['index']]])
assert trusted_cop_address_result.get('exit_code', -1) == 0, "Trusted cop address error"
trusted_cop_address = Cell.one_from_boc(b64decode(trusted_cop_address_result['stack'][0][1]['bytes'])).begin_parse().read_msg_addr().to_string(1, 1, 1)
@@ -72,7 +68,7 @@ class UserContentIndexationMixin:
self.owner_address = cc_indexator_data['owner_address']
self.type = 'nft/listen'
self.content_id = stored_content.id
db_session.commit()
await db_session.commit()
except BaseException as e:
errored = True
make_log("UserContent", f"Error: {e}" + '\n' + traceback.format_exc(), level="error")
@@ -80,7 +76,6 @@ class UserContentIndexationMixin:
if errored is True:
self.type = 'nft/unknown'
self.content_id = None
db_session.commit()
await db_session.commit()
+6 -6
View File
@@ -3,6 +3,7 @@ from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime
from sqlalchemy.orm import relationship
from app.core.models.base import AlchemyBase
from app.core.models.content.indexation_mixins import UserContentIndexationMixin
from datetime import datetime
class UserContent(AlchemyBase, UserContentIndexationMixin):
@@ -14,12 +15,12 @@ class UserContent(AlchemyBase, UserContentIndexationMixin):
owner_address = Column(String(1024), nullable=True)
code_hash = Column(String(128), nullable=True)
data_hash = Column(String(128), nullable=True)
updated = Column(DateTime, nullable=False, default=0)
updated = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
content_id = Column(Integer, ForeignKey('node_storage.id'), nullable=True)
created = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
meta = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default=dict)
user_id = Column(Integer, ForeignKey('users.id'), nullable=False)
wallet_connection_id = Column(Integer, ForeignKey('wallet_connections.id'), nullable=True)
status = Column(String(64), nullable=False, default='active') # 'transaction_requested'
@@ -41,9 +42,8 @@ class UserAction(AlchemyBase):
to_address = Column(String(1024), nullable=True)
from_address = Column(String(1024), nullable=True)
status = Column(String(128), nullable=True)
meta = Column(JSON, nullable=False, default={})
created = Column(DateTime, nullable=False, default=0)
meta = Column(JSON, nullable=False, default=dict)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
user = relationship('User', uselist=False, foreign_keys=[user_id])
content = relationship('StoredContent', uselist=False, foreign_keys=[content_id])
+3 -2
View File
@@ -1,5 +1,6 @@
from base58 import b58decode
from sqlalchemy import Column, Integer, String, DateTime, JSON
from datetime import datetime
from .base import AlchemyBase
@@ -15,12 +16,12 @@ class KnownKey(AlchemyBase):
public_key_hash = Column(String(64), nullable=False, unique=True) # base58
algo = Column(String(32), nullable=True, default=None)
meta = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default=dict)
# {
# "I_user_id": TRUSTED_USER_ID,
# }
created = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
# stored_content = relationship('StoredContent', back_populates='key')
+5 -6
View File
@@ -12,9 +12,9 @@ class KnownNode(AlchemyBase):
public_key = Column(String(256), nullable=False)
codebase_hash = Column(String(512), nullable=True) # Node software version
reputation = Column(Integer, nullable=False, default=0)
last_sync = Column(DateTime, nullable=False, default=datetime.now)
meta = Column(JSON, nullable=False, default={})
located_at = Column(DateTime, nullable=False, default=datetime.now)
last_sync = Column(DateTime, nullable=False, default=datetime.utcnow)
meta = Column(JSON, nullable=False, default=dict)
located_at = Column(DateTime, nullable=False, default=datetime.utcnow)
class KnownNodeIncident(AlchemyBase):
@@ -28,7 +28,7 @@ class KnownNodeIncident(AlchemyBase):
severity = Column(Integer, nullable=False, default=1) # Severity level (1-low to 5-critical)
resolved = Column(Boolean, nullable=False, default=False) # Whether the incident has been resolved
resolved_at = Column(DateTime, nullable=True) # Timestamp when the incident was resolved
meta = Column(JSON, nullable=False, default={}) # Additional metadata if needed
meta = Column(JSON, nullable=False, default=dict) # Additional metadata if needed
class RemoteContentIndex(AlchemyBase):
@@ -41,7 +41,6 @@ class RemoteContentIndex(AlchemyBase):
decrypted_hash = Column(String(128), nullable=True) # Decrypted content hash, available once permission is granted
ton_address = Column(String(128), nullable=True) # TON network address for the content
onchain_index = Column(Integer, nullable=True) # Onchain index or reference on a blockchain
meta = Column(JSON, nullable=False, default={}) # Additional metadata for flexible content description
meta = Column(JSON, nullable=False, default=dict) # Additional metadata for flexible content description
last_updated = Column(DateTime, nullable=False, default=datetime.utcnow) # Timestamp of the last update
created_at = Column(DateTime, nullable=False, default=datetime.utcnow) # Record creation timestamp
+53 -4
View File
@@ -25,7 +25,8 @@ class StoredContent(AlchemyBase, AudioContentMixin):
status = Column(String(32), nullable=True)
filename = Column(String(1024), nullable=False)
meta = Column(JSON, nullable=False, default={})
# Use a factory for JSON default to avoid shared mutable dict
meta = Column(JSON, nullable=False, default=dict)
user_id = Column(Integer, ForeignKey('users.id'), nullable=True)
owner_address = Column(String(1024), nullable=True)
@@ -35,9 +36,11 @@ class StoredContent(AlchemyBase, AudioContentMixin):
telegram_cid = Column(String(1024), nullable=True)
codebase_version = Column(Integer, nullable=True)
created = Column(DateTime, nullable=False, default=0)
updated = Column(DateTime, nullable=False, default=0)
disabled = Column(DateTime, nullable=False, default=0)
# Use proper datetime defaults; updated also auto-updates on change
created = Column(DateTime, nullable=False, default=datetime.utcnow)
updated = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
# Timestamp of when content was disabled; None means active
disabled = Column(DateTime, nullable=True, default=None)
disabled_by = Column(Integer, ForeignKey('users.id'), nullable=True, default=None)
encrypted = Column(Boolean, nullable=False, default=False)
@@ -96,6 +99,30 @@ class StoredContent(AlchemyBase, AudioContentMixin):
make_log("NodeStorage.open_content", f"Can't open content: {self.id} {e}", level='warning')
raise e
async def open_content_async(self, db_session, content_type=None):
from sqlalchemy import select
try:
decrypted_content = self if not self.encrypted else None
encrypted_content = self if self.encrypted else None
if not decrypted_content:
decrypted_content = (await db_session.execute(select(StoredContent).where(StoredContent.id == self.decrypted_content_id))).scalars().first()
else:
encrypted_content = (await db_session.execute(select(StoredContent).where(StoredContent.decrypted_content_id == self.id))).scalars().first()
assert decrypted_content, "Can't get decrypted content"
assert encrypted_content, "Can't get encrypted content"
_ct = content_type or decrypted_content.json_format()['content_type']
content_type = _ct.split('/')[0] if _ct else 'application'
return {
'encrypted_content': encrypted_content,
'decrypted_content': decrypted_content,
'content_type': content_type or 'application/x-binary'
}
except BaseException as e:
make_log("NodeStorage.open_content_async", f"Can't open content: {self.id} {e}", level='warning')
raise e
def json_format(self):
extra_fields = {}
if self.type.startswith('local'):
@@ -144,6 +171,16 @@ class StoredContent(AlchemyBase, AudioContentMixin):
with open(metadata_content.filepath, 'r') as f:
return json.loads(f.read())
async def metadata_json_async(self, db_session):
metadata_cid = self.meta.get('metadata_cid')
if not metadata_cid:
return None
metadata_content = await StoredContent.from_cid_async(db_session, metadata_cid)
import aiofiles
async with aiofiles.open(metadata_content.filepath, 'r') as f:
data = await f.read()
return json.loads(data)
@classmethod
def from_cid(cls, db_session, content_id):
if isinstance(content_id, str):
@@ -155,3 +192,15 @@ class StoredContent(AlchemyBase, AudioContentMixin):
assert content, "Content not found"
return content
@classmethod
async def from_cid_async(cls, db_session, content_id):
from sqlalchemy import select
if isinstance(content_id, str):
cid = ContentId.deserialize(content_id)
else:
cid = content_id
result = await db_session.execute(select(StoredContent).where(StoredContent.hash == cid.content_hash_b58))
content = result.scalars().first()
assert content, "Content not found"
return content
+16
View File
@@ -0,0 +1,16 @@
from .base import AlchemyBase
from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime, JSON, Boolean
from datetime import datetime
class PromoAction(AlchemyBase):
__tablename__ = 'promo_actions'
id = Column(Integer, autoincrement=True, primary_key=True)
user_id = Column(String(512), nullable=False)
user_internal_id = Column(Integer, ForeignKey('users.id'), nullable=True)
action_type = Column(String(64), nullable=False) # Type of action, e.g., 'referral', 'discount'
action_ref = Column(String(512), nullable=False) # Reference to the action, e.g., promo code
created = Column(DateTime, nullable=False, default=datetime.utcnow)
+26
View File
@@ -0,0 +1,26 @@
from .base import AlchemyBase
from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime, JSON, Boolean
from datetime import datetime
class BlockchainTask(AlchemyBase):
__tablename__ = 'blockchain_tasks'
id = Column(Integer, autoincrement=True, primary_key=True)
destination = Column(String(1024), nullable=False)
amount = Column(String(256), nullable=False)
payload = Column(String(4096), nullable=False)
epoch = Column(Integer, nullable=True)
seqno = Column(Integer, nullable=True)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
updated = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
user_id = Column(Integer, ForeignKey('users.id'), nullable=True)
meta = Column(JSON, nullable=False, default=dict)
status = Column(String(256), nullable=False)
transaction_hash = Column(String(1024), nullable=True)
transaction_lt = Column(String(1024), nullable=True)
+3 -3
View File
@@ -13,8 +13,8 @@ class UserBalance(AlchemyBase):
asset_id = Column(Integer, ForeignKey('assets.id'), nullable=False)
balance = Column(Float, nullable=False, default=0)
updated = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=0)
updated = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
user = relationship('User', uselist=False, foreign_keys=[user_id], back_populates='balances')
asset = relationship('Asset', uselist=False, foreign_keys=[asset_id], back_populates='balances')
@@ -32,7 +32,7 @@ class InternalTransaction(AlchemyBase):
spent_transaction_id = Column(Integer, ForeignKey('internal_transactions.id'), nullable=True)
type = Column(String(256), nullable=False, default="NOT_SPECIFIED")
created = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
user = relationship('User', uselist=False, back_populates='internal_transactions', foreign_keys=[user_id])
asset = relationship('Asset', uselist=False, foreign_keys=[asset_id])
+5 -4
View File
@@ -1,3 +1,4 @@
from datetime import datetime
from sqlalchemy import Column, Integer, String, BigInteger, DateTime, JSON
from sqlalchemy.orm import relationship
@@ -17,10 +18,11 @@ class User(AlchemyBase, DisplayMixin, TranslationCore, AuthenticationMixin_V1, W
username = Column(String(512), nullable=True)
lang_code = Column(String(8), nullable=False, default="en")
meta = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default=dict)
last_use = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=0)
last_use = Column(DateTime, nullable=False, default=datetime.utcnow)
updated = Column(DateTime, nullable=False, default=datetime.utcnow)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
balances = relationship('UserBalance', back_populates='user')
internal_transactions = relationship('InternalTransaction', back_populates='user')
@@ -30,4 +32,3 @@ class User(AlchemyBase, DisplayMixin, TranslationCore, AuthenticationMixin_V1, W
def __str__(self):
return f"User, {self.id}_{self.telegram_id} | Username: {self.username} " + '\\'
+24 -23
View File
@@ -10,18 +10,21 @@ from httpx import AsyncClient
class WalletMixin:
def wallet_connection(self, db_session):
return db_session.query(WalletConnection).filter(
WalletConnection.user_id == self.id,
WalletConnection.invalidated == False
).order_by(WalletConnection.created.desc()).first()
async def wallet_connection_async(self, db_session):
from sqlalchemy import select, and_, desc
result = await db_session.execute(
select(WalletConnection)
.where(and_(WalletConnection.user_id == self.id, WalletConnection.invalidated == False))
.order_by(WalletConnection.created.desc())
)
return result.scalars().first()
def wallet_address(self, db_session):
wallet_connection = self.wallet_connection(db_session)
return wallet_connection.wallet_address if wallet_connection else None
async def wallet_address_async(self, db_session):
wc = await self.wallet_connection_async(db_session)
return wc.wallet_address if wc else None
async def scan_owned_user_content(self, db_session):
user_wallet_address = self.wallet_address(db_session)
user_wallet_address = await self.wallet_address_async(db_session)
async def get_nft_items_list():
try:
@@ -40,9 +43,8 @@ class WalletMixin:
item_address = Address(nft_item['address']).to_string(1, 1, 1)
owner_address = Address(nft_item['owner']['address']).to_string(1, 1, 1)
user_content = db_session.query(UserContent).filter(
UserContent.onchain_address == item_address
).first()
from sqlalchemy import select
user_content = (await db_session.execute(select(UserContent).where(UserContent.onchain_address == item_address))).scalars().first()
if user_content:
continue
@@ -57,18 +59,18 @@ class WalletMixin:
created=datetime.now(),
meta={},
user_id=self.id,
wallet_connection_id=self.wallet_connection(db_session).id,
wallet_connection_id=(await self.wallet_connection_async(db_session)).id,
status="active"
)
db_session.add(user_content)
db_session.commit()
await db_session.commit()
make_log(self, f"New onchain NFT found: {item_address}", level='info')
async def ____scan_owned_user_content(self, db_session):
page_id = -1
page_size = 100
have_next_page = True
user_wallet_address = self.wallet_address(db_session)
user_wallet_address = await self.wallet_address_async(db_session)
while have_next_page:
page_id += 1
nfts_list = await toncenter.get_nft_items(limit=100, offset=page_id * page_size, owner_address=user_wallet_address)
@@ -81,9 +83,8 @@ class WalletMixin:
item_address = Address(nft_item['address']).to_string(1, 1, 1)
owner_address = Address(nft_item['owner_address']).to_string(1, 1, 1)
user_content = db_session.query(UserContent).filter(
UserContent.onchain_address == item_address
).first()
from sqlalchemy import select
user_content = (await db_session.execute(select(UserContent).where(UserContent.onchain_address == item_address))).scalars().first()
if user_content:
continue
@@ -105,11 +106,11 @@ class WalletMixin:
'metadata_uri': nft_content,
},
user_id=self.id,
wallet_connection_id=self.wallet_connection(db_session).id,
wallet_connection_id=(await self.wallet_connection_async(db_session)).id,
status="active"
)
db_session.add(user_content)
db_session.commit()
await db_session.commit()
make_log(self, f"New onchain NFT found: {item_address}", level='info')
except BaseException as e:
@@ -122,6 +123,6 @@ class WalletMixin:
except BaseException as e:
make_log(self, f"Error while scanning user content: {e}", level='error')
return self.db_session.query(UserContent).filter(
UserContent.user_id == self.id
).offset(offset).limit(limit).all()
from sqlalchemy import select
result = await db_session.execute(select(UserContent).where(UserContent.user_id == self.id).offset(offset).limit(limit))
return result.scalars().all()
+3 -2
View File
@@ -2,6 +2,7 @@
from sqlalchemy import Column, BigInteger, Integer, String, ForeignKey, DateTime, JSON, Boolean
from sqlalchemy.orm import relationship
from .base import AlchemyBase
from datetime import datetime
class UserActivity(AlchemyBase):
@@ -9,10 +10,10 @@ class UserActivity(AlchemyBase):
id = Column(Integer, autoincrement=True, primary_key=True)
type = Column(String(64), nullable=False)
meta = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default=dict)
user_id = Column(Integer, ForeignKey('users.id'), nullable=True)
user_ip = Column(String(64), nullable=True)
created = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
user = relationship('User', uselist=False, foreign_keys=[user_id])
+5 -4
View File
@@ -1,5 +1,6 @@
from sqlalchemy import Column, Integer, String, ForeignKey, DateTime, JSON, Boolean
from sqlalchemy.orm import relationship
from datetime import datetime
from .base import AlchemyBase
@@ -15,11 +16,11 @@ class WalletConnection(AlchemyBase):
wallet_address = Column(String(1024), nullable=False)
keys = Column(JSON, nullable=False, default={})
meta = Column(JSON, nullable=False, default={})
keys = Column(JSON, nullable=False, default=dict)
meta = Column(JSON, nullable=False, default=dict)
created = Column(DateTime, nullable=False, default=0)
updated = Column(DateTime, nullable=False, default=0)
created = Column(DateTime, nullable=False, default=datetime.utcnow)
updated = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
invalidated = Column(Boolean, nullable=False, default=True)
without_pk = Column(Boolean, nullable=False, default=False)
+1 -1
View File
@@ -68,7 +68,7 @@ file_handler.setLevel(logging.DEBUG)
file_handler.setFormatter(logging.Formatter(FORMAT_STRING))
logger.addHandler(file_handler)
if os.getenv('APP_ENABLE_STDOUT_LOGS', '0') == '1':
if int(os.getenv('APP_ENABLE_STDOUT_LOGS', '0')) == 1:
stdout_handler = logging.StreamHandler()
stdout_handler.setLevel(logging.DEBUG)
stdout_handler.setFormatter(logging.Formatter(FORMAT_STRING))
+41 -29
View File
@@ -1,45 +1,57 @@
import time
from contextlib import contextmanager
from contextlib import asynccontextmanager
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from sqlalchemy.sql import text
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from app.core._config import MYSQL_URI, MYSQL_DATABASE
from app.core._config import DATABASE_URL
from app.core.logger import make_log
from sqlalchemy.pool import NullPool
engine = create_engine(MYSQL_URI, poolclass=NullPool) #, echo=True)
Session = sessionmaker(bind=engine)
database_initialized = False
while not database_initialized:
def _to_async_dsn(url: str) -> str:
# Convert psycopg2 DSN to asyncpg DSN
# postgresql+psycopg2://user:pass@host:5432/db -> postgresql+asyncpg://user:pass@host:5432/db
return url.replace("+psycopg2", "+asyncpg")
# Async engine for PostgreSQL
engine = create_async_engine(
_to_async_dsn(DATABASE_URL),
pool_size=10,
max_overflow=20,
pool_timeout=30,
pool_recycle=1800,
pool_pre_ping=True,
)
AsyncSessionLocal = async_sessionmaker(engine, expire_on_commit=False, class_=AsyncSession)
async def wait_db_ready():
ready = False
while not ready:
try:
with Session() as session:
databases_list = session.execute(text("SHOW DATABASES;"))
databases_list = [row[0] for row in databases_list]
make_log("SQL", 'Database list: ' + str(databases_list), level='debug')
assert MYSQL_DATABASE in databases_list, 'Database not found'
database_initialized = True
async with engine.connect() as conn:
await conn.execute(text("SELECT 1"))
ready = True
except Exception as e:
make_log("SQL", 'MariaDB is not ready yet: ' + str(e), level='debug')
make_log("SQL", 'PostgreSQL is not ready yet: ' + str(e), level='debug')
time.sleep(1)
engine = create_engine(f"{MYSQL_URI}/{MYSQL_DATABASE}", poolclass=NullPool)
Session = sessionmaker(bind=engine)
@contextmanager
def db_session(auto_commit=False):
_session = Session()
@asynccontextmanager
async def db_session(auto_commit: bool = False):
session: AsyncSession = AsyncSessionLocal()
try:
yield _session
if auto_commit is True:
_session.commit()
yield session
if auto_commit:
await session.commit()
except BaseException as e:
_session.rollback()
await session.rollback()
raise e
finally:
_session.close()
await session.close()
def new_session() -> AsyncSession:
return AsyncSessionLocal()
+15 -18
View File
@@ -1,19 +1,19 @@
from datetime import datetime
from sqlalchemy import select, and_, func
from app.core.logger import make_log
from app.core.models import Memory, User, UserBalance, Asset, InternalTransaction
from app.core.storage import db_session
def get_user_balance(session, user: User, asset: Asset) -> UserBalance:
async def get_user_balance(session, user: User, asset: Asset) -> UserBalance:
assert user, "No user"
assert asset, "No asset"
result = session.query(UserBalance).filter(
UserBalance.user_id == user.id,
UserBalance.asset_id == asset.id
)
results_count = result.count()
if results_count == 0:
result = await session.execute(select(UserBalance).where(
and_(UserBalance.user_id == user.id, UserBalance.asset_id == asset.id)
))
row = result.scalars().first()
if not row:
user_balance = UserBalance(
user_id=user.id,
asset_id=asset.id,
@@ -21,12 +21,9 @@ def get_user_balance(session, user: User, asset: Asset) -> UserBalance:
created=datetime.now(),
)
session.add(user_balance)
session.commit()
return get_user_balance(session, user, asset)
elif results_count == 1:
return result.first()
else:
raise Exception(f"Multiple user balances found: {results_count}")
await session.commit()
return await get_user_balance(session, user, asset)
return row
async def make_internal_transaction(
@@ -46,13 +43,13 @@ async def make_internal_transaction(
raise Exception(f"Invalid amount: {amount}")
abs_amount = abs(amount)
with db_session(auto_commit=False) as session:
async with db_session(auto_commit=False) as session:
async with memory.transaction():
user = session.query(User).filter_by(id=user_id).first()
user = (await session.execute(select(User).where(User.id == user_id))).scalars().first()
assert user, "No user"
asset = session.query(Asset).filter_by(id=asset_id).first()
asset = (await session.execute(select(Asset).where(Asset.id == asset_id))).scalars().first()
assert asset, "No asset"
user_balance = get_user_balance(session, user, asset)
user_balance = await get_user_balance(session, user, asset)
assert user_balance, "No user balance"
if is_spent is True:
if abs_amount > user_balance.balance:
@@ -71,6 +68,6 @@ async def make_internal_transaction(
created=datetime.now(),
)
session.add(internal_transaction)
session.commit()
await session.commit()
make_log(user, f"Made internal transaction: {'-' if is_spent else ''}{abs_amount} {asset.symbol}, type: {type}")
-99
View File
@@ -1,99 +0,0 @@
version: '3'
services:
maria_db:
image: mariadb:11.2
ports:
- "3307:3306"
env_file:
- .env
volumes:
- /Storage/sqlStorage:/var/lib/mysql
healthcheck:
test: [ "CMD", "healthcheck.sh", "--connect", "--innodb_initialized" ]
interval: 10s
timeout: 5s
retries: 3
app:
build:
context: .
dockerfile: Dockerfile
command: python -m app
env_file:
- .env
links:
- maria_db
ports:
- "15100:15100"
volumes:
- /Storage/logs:/app/logs
- /Storage/storedContent:/app/data
depends_on:
maria_db:
condition: service_healthy
indexer: # Отправка уведомления о появлении новой NFT-listen. Установка CID поля у всего контента. Проверка следующего за последним индексом item коллекции и поиск нового контента, отправка информации о том что контент найден его загружателю. Присваивание encrypted_content onchain_index
build:
context: .
dockerfile: Dockerfile
command: python -m app indexer
env_file:
- .env
links:
- maria_db
volumes:
- /Storage/logs:/app/logs
- /Storage/storedContent:/app/data
depends_on:
maria_db:
condition: service_healthy
ton_daemon: # Работа с TON-сетью. Задачи сервисного кошелька и деплой контрактов
build:
context: .
dockerfile: Dockerfile
command: python -m app ton_daemon
env_file:
- .env
links:
- maria_db
volumes:
- /Storage/logs:/app/logs
- /Storage/storedContent:/app/data
depends_on:
maria_db:
condition: service_healthy
license_index: # Проверка кошельков пользователей на новые NFT. Опрос этих NFT на определяемый GET-метод по которому мы определяем что это определенная лицензия и сохранение информации по ней
build:
context: .
dockerfile: Dockerfile
command: python -m app license_index
env_file:
- .env
links:
- maria_db
volumes:
- /Storage/logs:/app/logs
- /Storage/storedContent:/app/data
depends_on:
maria_db:
condition: service_healthy
convert_process:
build:
context: .
dockerfile: Dockerfile
command: python -m app convert_process
env_file:
- .env
links:
- maria_db
volumes:
- /Storage/logs:/app/logs
- /Storage/storedContent:/app/data
- /var/run/docker.sock:/var/run/docker.sock
depends_on:
maria_db:
condition: service_healthy
Binary file not shown.
+14 -2
View File
@@ -148,8 +148,11 @@ msgid "buyTrackListenLicense_button"
msgstr "💶 Купить лицензию ({price} TON)"
#: app/core/models/_telegram/templates/player.py:117
msgid "p_playerContext_preview"
msgstr "<i>Это демонстрационная версия трека. Чтобы прослушать полную версию, необходимо приобрести лицензию 🎧.</i>"
msgid "p_playerAudioContext_preview"
msgstr "<i>Превью аудиозаписи. Полная версия доступна только по лицензии 🎧</i>"
msgid "p_playerVideoContext_preview"
msgstr "<i>Превью видео. Полная версия доступна только по лицензии 🎬</i>"
#: app/client_bot/routers/content.py:32
msgid "error_contentNotFound"
@@ -194,3 +197,12 @@ msgstr "🔗 Поделиться"
msgid "p_shareLinkContext"
msgstr "🎉 Наслаждайтесь {title} на MY!"
msgid "p_uploadContentTxPromo"
msgstr ""
"🎉 Вам доступно ещё <b>{free_count}</b> бесплатных приветственных загрузок контента! "
"Контент <b>{title}</b> уже находится в процессе загрузки. Как только блокчейн обработает транзакцию, "
"вы получите NFT-лицензию."
msgid "p_playerContext_contentNotReady"
msgstr "⚠️ Контент, который вы хотите просмотреть, ещё не готов. Пожалуйста, попробуйте позже."
+4 -3
View File
@@ -2,11 +2,12 @@ sanic==21.9.1
websockets==10.0
sqlalchemy==2.0.23
python-dotenv==1.0.0
pymysql==1.1.0
psycopg2-binary==2.9.9
asyncpg==0.29.0
aiogram==3.13.0
pytonconnect==0.3.0
base58==2.1.1
tonsdk==1.0.13
git+https://github.com/tonfactory/tonsdk.git@3ebbf0b702f48c2519e4c6c425f9514f673b9d48#egg=tonsdk
httpx==0.25.0
docker==7.0.0
pycryptodome==3.20.0
@@ -15,4 +16,4 @@ aiofiles==23.2.1
pydub==0.25.1
pillow==10.2.0
ffmpeg-python==0.2.0
python-magic==0.4.27