smashed updated
This commit is contained in:
1 parent
77921ba6a8
commit
1da0b26320
14 files changed
+321
-96
No files matched your search
@@ -8,7 +8,7 @@ from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import List, Optional, Tuple
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy import select, and_, or_
|
||||
|
||||
from app.core.logger import make_log
|
||||
from app.core.storage import db_session
|
||||
@@ -196,17 +196,45 @@ async def _convert_content(ec: EncryptedContent, staging: PlainStaging):
|
||||
plain_filename = f"{ec.encrypted_cid}.{input_ext}" if input_ext else ec.encrypted_cid
|
||||
async with db_session() as session:
|
||||
existing = (await session.execute(select(StoredContent).where(StoredContent.hash == file_hash))).scalars().first()
|
||||
if not existing:
|
||||
if existing:
|
||||
sc = existing
|
||||
sc.type = sc.type or "local/content_bin"
|
||||
sc.filename = plain_filename
|
||||
sc.meta = {
|
||||
**(sc.meta or {}),
|
||||
'encrypted_cid': ec.encrypted_cid,
|
||||
'kind': 'original',
|
||||
'content_type': ec.content_type,
|
||||
}
|
||||
sc.updated = datetime.utcnow()
|
||||
else:
|
||||
sc = StoredContent(
|
||||
type="local/content_bin",
|
||||
hash=file_hash,
|
||||
user_id=None,
|
||||
filename=plain_filename,
|
||||
meta={'encrypted_cid': ec.encrypted_cid, 'kind': 'original'},
|
||||
meta={
|
||||
'encrypted_cid': ec.encrypted_cid,
|
||||
'kind': 'original',
|
||||
'content_type': ec.content_type,
|
||||
},
|
||||
created=datetime.utcnow(),
|
||||
)
|
||||
session.add(sc)
|
||||
await session.flush()
|
||||
|
||||
encrypted_records = (await session.execute(select(StoredContent).where(StoredContent.hash == encrypted_hash_b58))).scalars().all()
|
||||
for encrypted_sc in encrypted_records:
|
||||
meta = dict(encrypted_sc.meta or {})
|
||||
converted = dict(meta.get('converted_content') or {})
|
||||
converted['original'] = file_hash
|
||||
meta['converted_content'] = converted
|
||||
if 'content_type' not in meta:
|
||||
meta['content_type'] = ec.content_type
|
||||
encrypted_sc.meta = meta
|
||||
encrypted_sc.decrypted_content_id = sc.id
|
||||
encrypted_sc.updated = datetime.utcnow()
|
||||
|
||||
derivative = ContentDerivative(
|
||||
content_id=ec.id,
|
||||
kind='decrypted_original',
|
||||
@@ -341,10 +369,17 @@ async def _convert_content(ec: EncryptedContent, staging: PlainStaging):
|
||||
|
||||
async def _pick_pending(limit: int) -> List[Tuple[EncryptedContent, PlainStaging]]:
|
||||
async with db_session() as session:
|
||||
# Find A/V contents with preview_enabled and no ready low/low_preview derivatives yet
|
||||
ecs = (await session.execute(select(EncryptedContent).where(
|
||||
EncryptedContent.preview_enabled == True
|
||||
).order_by(EncryptedContent.created_at.desc()))).scalars().all()
|
||||
# Include preview-enabled media and non-media content that need decrypted originals
|
||||
non_media_filter = and_(
|
||||
EncryptedContent.content_type.isnot(None),
|
||||
~EncryptedContent.content_type.like('audio/%'),
|
||||
~EncryptedContent.content_type.like('video/%'),
|
||||
)
|
||||
ecs = (await session.execute(
|
||||
select(EncryptedContent)
|
||||
.where(or_(EncryptedContent.preview_enabled == True, non_media_filter))
|
||||
.order_by(EncryptedContent.created_at.desc())
|
||||
)).scalars().all()
|
||||
|
||||
picked: List[Tuple[EncryptedContent, PlainStaging]] = []
|
||||
for ec in ecs:
|
||||
@@ -365,7 +400,12 @@ async def _pick_pending(limit: int) -> List[Tuple[EncryptedContent, PlainStaging
|
||||
# Check if derivatives already ready
|
||||
rows = (await session.execute(select(ContentDerivative).where(ContentDerivative.content_id == ec.id))).scalars().all()
|
||||
kinds_ready = {r.kind for r in rows if r.status == 'ready'}
|
||||
required = {'decrypted_low', 'decrypted_high'} if ec.content_type.startswith('audio/') else {'decrypted_low', 'decrypted_high', 'decrypted_preview'}
|
||||
if ec.content_type.startswith('audio/'):
|
||||
required = {'decrypted_low', 'decrypted_high'}
|
||||
elif ec.content_type.startswith('video/'):
|
||||
required = {'decrypted_low', 'decrypted_high', 'decrypted_preview'}
|
||||
else:
|
||||
required = {'decrypted_original'}
|
||||
if required.issubset(kinds_ready):
|
||||
continue
|
||||
# Always decrypt from IPFS using local or remote key
|
||||
|
||||
@@ -86,11 +86,14 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
|
||||
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._telegram_bot, chat_id=wallet_owner_user.telegram_id, user=wallet_owner_user, db_session=session)
|
||||
meta_title = content_metadata.get('title') or content_metadata.get('name') or 'Unknown'
|
||||
meta_artist = content_metadata.get('artist')
|
||||
formatted_title = f"{meta_artist} – {meta_title}" if meta_artist else meta_title
|
||||
await wallet_owner_bot.send_message(
|
||||
user.translated('p_licenseWasBought').format(
|
||||
username=user.front_format(),
|
||||
nft_address=f'"https://tonviewer.com/{new_license.onchain_address}"',
|
||||
content_title=content_metadata.get('name', 'Unknown'),
|
||||
content_title=formatted_title,
|
||||
),
|
||||
message_type='notification',
|
||||
)
|
||||
|
||||
Reference in new issue
Block a user