events & global sync. unstable
This commit is contained in:
1 parent
77921ba6a8
commit
f140181c45
21 files changed
+1196
-31
No files matched your search
@@ -1,4 +1,6 @@
|
||||
from sqlalchemy.ext.asyncio import AsyncEngine
|
||||
from sqlalchemy import text
|
||||
|
||||
from app.core.models import BlockchainTask
|
||||
from app.core.models.base import AlchemyBase
|
||||
|
||||
@@ -9,4 +11,36 @@ async def create_db_tables(engine: AsyncEngine):
|
||||
BlockchainTask()
|
||||
async with engine.begin() as conn:
|
||||
await conn.run_sync(AlchemyBase.metadata.create_all)
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE users
|
||||
ADD COLUMN IF NOT EXISTS is_admin BOOLEAN DEFAULT FALSE
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS telegram_id BIGINT
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS paid_at TIMESTAMPTZ
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS payment_tx_id VARCHAR(256)
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS payment_node_id VARCHAR(128)
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS payment_node_public_host VARCHAR(256)
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS bot_username VARCHAR(128)
|
||||
"""))
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE stars_invoices
|
||||
ADD COLUMN IF NOT EXISTS is_remote BOOLEAN DEFAULT FALSE
|
||||
"""))
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
import asyncio
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
from urllib.parse import urlencode
|
||||
|
||||
import httpx
|
||||
from sqlalchemy import select
|
||||
|
||||
from app.core.logger import make_log
|
||||
from app.core.storage import db_session
|
||||
from app.core.models import KnownNode, NodeEvent
|
||||
from app.core.events.service import (
|
||||
store_remote_events,
|
||||
upsert_cursor,
|
||||
LOCAL_PUBLIC_KEY,
|
||||
)
|
||||
from app.core.models.events import NodeEventCursor
|
||||
from app.core._secrets import hot_pubkey, hot_seed
|
||||
from app.core.network.nodesig import sign_headers
|
||||
from base58 import b58encode
|
||||
|
||||
|
||||
def _node_public_base(node: KnownNode) -> Optional[str]:
|
||||
meta = node.meta or {}
|
||||
public_host = (meta.get('public_host') or '').strip()
|
||||
if public_host:
|
||||
base = public_host.rstrip('/')
|
||||
if base.startswith('http://') or base.startswith('https://'):
|
||||
return base
|
||||
scheme = 'https' if node.port == 443 else 'http'
|
||||
return f"{scheme}://{base.lstrip('/')}"
|
||||
scheme = 'https' if node.port == 443 else 'http'
|
||||
host = (node.ip or '').strip()
|
||||
if not host:
|
||||
return None
|
||||
default_port = 443 if scheme == 'https' else 80
|
||||
if node.port and node.port != default_port:
|
||||
return f"{scheme}://{host}:{node.port}"
|
||||
return f"{scheme}://{host}"
|
||||
|
||||
|
||||
async def _fetch_events_for_node(node: KnownNode, limit: int = 100) -> Tuple[List[Dict], int]:
|
||||
base = _node_public_base(node)
|
||||
if not base:
|
||||
return [], 0
|
||||
async with db_session() as session:
|
||||
cursor = (await session.execute(
|
||||
select(NodeEventCursor).where(NodeEventCursor.source_public_key == node.public_key)
|
||||
)).scalar_one_or_none()
|
||||
since = cursor.last_seq if cursor else 0
|
||||
query = urlencode({"since": since, "limit": limit})
|
||||
path = f"/api/v1/network.events?{query}"
|
||||
url = f"{base}{path}"
|
||||
pk_b58 = b58encode(hot_pubkey).decode()
|
||||
headers = sign_headers("GET", path, b"", hot_seed, pk_b58)
|
||||
async with httpx.AsyncClient(timeout=20.0) as client:
|
||||
try:
|
||||
resp = await client.get(url, headers=headers)
|
||||
if resp.status_code == 403:
|
||||
make_log("Events", f"Access denied by node {node.public_key}", level="warning")
|
||||
return [], since
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
except Exception as exc:
|
||||
make_log("Events", f"Fetch events failed from {node.public_key}: {exc}", level="debug")
|
||||
return [], since
|
||||
events = data.get("events") or []
|
||||
next_since = int(data.get("next_since") or since)
|
||||
return events, next_since
|
||||
|
||||
|
||||
async def _apply_event(session, event: NodeEvent):
|
||||
if event.event_type == "stars_payment":
|
||||
from app.core.models import StarsInvoice
|
||||
payload = event.payload or {}
|
||||
invoice_id = payload.get("invoice_id")
|
||||
telegram_id = payload.get("telegram_id")
|
||||
content_hash = payload.get("content_hash")
|
||||
amount = payload.get("amount")
|
||||
if not invoice_id or not telegram_id or not content_hash:
|
||||
return
|
||||
invoice = (await session.execute(select(StarsInvoice).where(StarsInvoice.external_id == invoice_id))).scalar_one_or_none()
|
||||
if not invoice:
|
||||
invoice = StarsInvoice(
|
||||
external_id=invoice_id,
|
||||
user_id=payload.get("user_id"),
|
||||
type=payload.get('type') or 'access',
|
||||
telegram_id=telegram_id,
|
||||
amount=amount,
|
||||
content_hash=content_hash,
|
||||
paid=True,
|
||||
paid_at=event.created_at,
|
||||
payment_node_id=payload.get("payment_node", {}).get("public_key"),
|
||||
payment_node_public_host=payload.get("payment_node", {}).get("public_host"),
|
||||
bot_username=payload.get("bot_username"),
|
||||
is_remote=True,
|
||||
)
|
||||
session.add(invoice)
|
||||
else:
|
||||
invoice.paid = True
|
||||
invoice.paid_at = invoice.paid_at or event.created_at
|
||||
invoice.payment_node_id = payload.get("payment_node", {}).get("public_key")
|
||||
invoice.payment_node_public_host = payload.get("payment_node", {}).get("public_host")
|
||||
invoice.bot_username = payload.get("bot_username") or invoice.bot_username
|
||||
invoice.telegram_id = telegram_id or invoice.telegram_id
|
||||
invoice.is_remote = invoice.is_remote or True
|
||||
if payload.get('type'):
|
||||
invoice.type = payload['type']
|
||||
event.status = 'applied'
|
||||
event.applied_at = event.applied_at or event.received_at
|
||||
elif event.event_type == "content_indexed":
|
||||
# The index scout will pick up via remote_content_index; we only mark event applied
|
||||
event.status = 'recorded'
|
||||
elif event.event_type == "node_registered":
|
||||
event.status = 'recorded'
|
||||
else:
|
||||
event.status = 'recorded'
|
||||
|
||||
|
||||
async def main_fn(memory):
|
||||
make_log("Events", "Sync service started", level="info")
|
||||
while True:
|
||||
try:
|
||||
async with db_session() as session:
|
||||
nodes = (await session.execute(select(KnownNode))).scalars().all()
|
||||
trusted_nodes = [
|
||||
n for n in nodes
|
||||
if isinstance(n.meta, dict) and n.meta.get("role") == "trusted" and n.public_key != LOCAL_PUBLIC_KEY
|
||||
]
|
||||
trusted_keys = {n.public_key for n in trusted_nodes}
|
||||
for node in trusted_nodes:
|
||||
events, next_since = await _fetch_events_for_node(node)
|
||||
if not events:
|
||||
if next_since:
|
||||
async with db_session() as session:
|
||||
await upsert_cursor(session, node.public_key, next_since, node.meta.get("public_host") if isinstance(node.meta, dict) else None)
|
||||
await session.commit()
|
||||
continue
|
||||
async with db_session() as session:
|
||||
stored = await store_remote_events(
|
||||
session,
|
||||
events,
|
||||
allowed_public_keys=trusted_keys,
|
||||
)
|
||||
for ev in stored:
|
||||
await _apply_event(session, ev)
|
||||
if stored:
|
||||
await session.commit()
|
||||
await upsert_cursor(session, node.public_key, next_since, node.meta.get("public_host") if isinstance(node.meta, dict) else None)
|
||||
await session.commit()
|
||||
except Exception as exc:
|
||||
make_log("Events", f"Sync loop error: {exc}", level="error")
|
||||
await asyncio.sleep(10)
|
||||
@@ -1,5 +1,6 @@
|
||||
import asyncio
|
||||
import os
|
||||
from datetime import datetime
|
||||
from typing import List, Optional
|
||||
|
||||
import httpx
|
||||
@@ -10,9 +11,11 @@ from sqlalchemy import select
|
||||
|
||||
from app.core.logger import make_log
|
||||
from app.core.storage import db_session
|
||||
from app.core.models.my_network import KnownNode
|
||||
from app.core.models.my_network import KnownNode, RemoteContentIndex
|
||||
from app.core.models.events import NodeEvent
|
||||
from app.core.models.content_v3 import EncryptedContent, ContentDerivative
|
||||
from app.core.ipfs_client import pin_add, pin_ls, find_providers, swarm_connect, add_streamed_file
|
||||
from app.core.events.service import LOCAL_PUBLIC_KEY
|
||||
|
||||
|
||||
INTERVAL_SEC = 60
|
||||
@@ -105,6 +108,71 @@ async def upsert_content(item: dict):
|
||||
make_log('index_scout_v3', f"thumbnail fetch failed for {cid}: {e}", level='warning')
|
||||
|
||||
|
||||
def _node_base_url(node: KnownNode) -> Optional[str]:
|
||||
meta = node.meta or {}
|
||||
public_host = (meta.get('public_host') or '').strip()
|
||||
if public_host:
|
||||
base = public_host.rstrip('/')
|
||||
if base.startswith('http://') or base.startswith('https://'):
|
||||
return base
|
||||
scheme = 'https' if node.port == 443 else 'http'
|
||||
return f"{scheme}://{base.lstrip('/')}"
|
||||
scheme = 'https' if node.port == 443 else 'http'
|
||||
host = (node.ip or '').strip()
|
||||
if not host:
|
||||
return None
|
||||
default_port = 443 if scheme == 'https' else 80
|
||||
if node.port and node.port != default_port:
|
||||
return f"{scheme}://{host}:{node.port}"
|
||||
return f"{scheme}://{host}"
|
||||
|
||||
|
||||
async def _update_remote_index(node_id: int, items: List[dict], *, incremental: bool):
|
||||
if not items:
|
||||
return
|
||||
async with db_session() as session:
|
||||
existing_rows = (await session.execute(
|
||||
select(RemoteContentIndex).where(RemoteContentIndex.remote_node_id == node_id)
|
||||
)).scalars().all()
|
||||
existing_map = {row.encrypted_hash: row for row in existing_rows if row.encrypted_hash}
|
||||
seen = set()
|
||||
now = datetime.utcnow()
|
||||
for item in items:
|
||||
cid = item.get('encrypted_cid')
|
||||
if not cid:
|
||||
continue
|
||||
seen.add(cid)
|
||||
payload_meta = {
|
||||
'title': item.get('title'),
|
||||
'description': item.get('description'),
|
||||
'size_bytes': item.get('size_bytes'),
|
||||
'preview_enabled': item.get('preview_enabled'),
|
||||
'preview_conf': item.get('preview_conf'),
|
||||
'issuer_node_id': item.get('issuer_node_id'),
|
||||
'salt_b64': item.get('salt_b64'),
|
||||
}
|
||||
meta_clean = {k: v for k, v in payload_meta.items() if v is not None}
|
||||
row = existing_map.get(cid)
|
||||
if row:
|
||||
row.content_type = item.get('content_type') or row.content_type
|
||||
row.meta = {**(row.meta or {}), **meta_clean}
|
||||
row.last_updated = now
|
||||
else:
|
||||
row = RemoteContentIndex(
|
||||
remote_node_id=node_id,
|
||||
content_type=item.get('content_type') or 'application/octet-stream',
|
||||
encrypted_hash=cid,
|
||||
meta=meta_clean,
|
||||
last_updated=now,
|
||||
)
|
||||
session.add(row)
|
||||
if not incremental and existing_map:
|
||||
for hash_value, row in list(existing_map.items()):
|
||||
if hash_value not in seen:
|
||||
await session.delete(row)
|
||||
await session.commit()
|
||||
|
||||
|
||||
async def main_fn(memory):
|
||||
make_log('index_scout_v3', 'Service started', level='info')
|
||||
sem = None
|
||||
@@ -119,8 +187,70 @@ async def main_fn(memory):
|
||||
sem = asyncio.Semaphore(max_pins)
|
||||
async with db_session() as session:
|
||||
nodes = (await session.execute(select(KnownNode))).scalars().all()
|
||||
node_by_pk = {n.public_key: n for n in nodes if n.public_key}
|
||||
async with db_session() as session:
|
||||
pending_events = (await session.execute(
|
||||
select(NodeEvent)
|
||||
.where(NodeEvent.event_type == 'content_indexed', NodeEvent.status.in_(('recorded', 'local', 'processing')))
|
||||
.order_by(NodeEvent.created_at.asc())
|
||||
.limit(25)
|
||||
)).scalars().all()
|
||||
for ev in pending_events:
|
||||
if ev.status != 'processing':
|
||||
ev.status = 'processing'
|
||||
await session.commit()
|
||||
for ev in pending_events:
|
||||
payload = ev.payload or {}
|
||||
cid = payload.get('encrypted_cid') or payload.get('content_cid')
|
||||
if ev.origin_public_key == LOCAL_PUBLIC_KEY:
|
||||
async with db_session() as session:
|
||||
ref = await session.get(NodeEvent, ev.id)
|
||||
if ref:
|
||||
ref.status = 'applied'
|
||||
ref.applied_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
continue
|
||||
if not cid:
|
||||
async with db_session() as session:
|
||||
ref = await session.get(NodeEvent, ev.id)
|
||||
if ref:
|
||||
ref.status = 'applied'
|
||||
ref.applied_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
continue
|
||||
node = node_by_pk.get(ev.origin_public_key)
|
||||
if not node:
|
||||
async with db_session() as session:
|
||||
node = (await session.execute(select(KnownNode).where(KnownNode.public_key == ev.origin_public_key))).scalar_one_or_none()
|
||||
if node:
|
||||
node_by_pk[node.public_key] = node
|
||||
if not node:
|
||||
make_log('index_scout_v3', f"Event {ev.uid} refers to unknown node {ev.origin_public_key}", level='debug')
|
||||
async with db_session() as session:
|
||||
ref = await session.get(NodeEvent, ev.id)
|
||||
if ref:
|
||||
ref.status = 'recorded'
|
||||
await session.commit()
|
||||
continue
|
||||
try:
|
||||
await _pin_one(node, cid)
|
||||
async with db_session() as session:
|
||||
ref = await session.get(NodeEvent, ev.id)
|
||||
if ref:
|
||||
ref.status = 'applied'
|
||||
ref.applied_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
except Exception as exc:
|
||||
make_log('index_scout_v3', f"Event pin failed for {cid}: {exc}", level='warning')
|
||||
async with db_session() as session:
|
||||
ref = await session.get(NodeEvent, ev.id)
|
||||
if ref:
|
||||
ref.status = 'recorded'
|
||||
await session.commit()
|
||||
for n in nodes:
|
||||
base = f"http://{n.ip}:{n.port}"
|
||||
base = _node_base_url(n)
|
||||
if not base:
|
||||
continue
|
||||
# jitter 0..30s per node to reduce stampede
|
||||
await asyncio.sleep(random.uniform(0, 30))
|
||||
etag = (n.meta or {}).get('index_etag')
|
||||
@@ -144,6 +274,10 @@ async def main_fn(memory):
|
||||
if not items:
|
||||
continue
|
||||
make_log('index_scout_v3', f"Fetched {len(items)} from {base}")
|
||||
try:
|
||||
await _update_remote_index(n.id, items, incremental=bool(since))
|
||||
except Exception as exc:
|
||||
make_log('index_scout_v3', f"remote index update failed for node {n.id}: {exc}", level='warning')
|
||||
|
||||
# Check disk watermark
|
||||
try:
|
||||
@@ -156,10 +290,10 @@ async def main_fn(memory):
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def _pin_one(cid: str):
|
||||
async def _pin_one(node: KnownNode, cid: str):
|
||||
async with sem:
|
||||
try:
|
||||
node_ipfs_meta = (n.meta or {}).get('ipfs') or {}
|
||||
node_ipfs_meta = (node.meta or {}).get('ipfs') or {}
|
||||
multiaddrs = node_ipfs_meta.get('multiaddrs') or []
|
||||
for addr in multiaddrs:
|
||||
try:
|
||||
@@ -196,11 +330,11 @@ async def main_fn(memory):
|
||||
except Exception as e:
|
||||
# Attempt HTTP gateway fallback before logging failure
|
||||
fallback_sources = []
|
||||
node_host = n.meta.get('public_host') if isinstance(n.meta, dict) else None
|
||||
node_host = node.meta.get('public_host') if isinstance(node.meta, dict) else None
|
||||
try:
|
||||
# Derive gateway host: prefer public_host domain if present
|
||||
parsed = urlparse(node_host) if node_host else None
|
||||
gateway_host = parsed.hostname if parsed and parsed.hostname else (n.ip or '').split(':')[0]
|
||||
gateway_host = parsed.hostname if parsed and parsed.hostname else (node.ip or '').split(':')[0]
|
||||
gateway_port = parsed.port if (parsed and parsed.port not in (None, 80, 443)) else 8080
|
||||
if gateway_host:
|
||||
gateway_url = f"http://{gateway_host}:{gateway_port}/ipfs/{cid}"
|
||||
@@ -234,7 +368,7 @@ async def main_fn(memory):
|
||||
cid = it.get('encrypted_cid')
|
||||
if cid:
|
||||
make_log('index_scout_v3', f"queue pin {cid}")
|
||||
tasks.append(asyncio.create_task(_pin_one(cid)))
|
||||
tasks.append(asyncio.create_task(_pin_one(n, cid)))
|
||||
if tasks:
|
||||
await asyncio.gather(*tasks)
|
||||
except Exception as e:
|
||||
|
||||
@@ -8,6 +8,7 @@ from sqlalchemy import String, and_, desc, cast
|
||||
from tonsdk.boc import Cell
|
||||
from tonsdk.utils import Address
|
||||
from app.core._config import CLIENT_TELEGRAM_BOT_USERNAME, PROJECT_HOST
|
||||
from app.core.events.service import record_event
|
||||
from app.core._blockchain.ton.platform import platform
|
||||
from app.core._blockchain.ton.toncenter import toncenter
|
||||
from app.core._utils.send_status import send_status
|
||||
@@ -284,6 +285,21 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
|
||||
**item_metadata_packed
|
||||
}
|
||||
encrypted_stored_content.content_id = item_content_cid_str
|
||||
try:
|
||||
await record_event(
|
||||
session,
|
||||
'content_indexed',
|
||||
{
|
||||
'onchain_index': item_index,
|
||||
'content_hash': item_content_hash_str,
|
||||
'encrypted_cid': item_content_cid_str,
|
||||
'item_address': item_address.to_string(1, 1, 1),
|
||||
'owner_address': item_owner_address.to_string(1, 1, 1) if item_owner_address else None,
|
||||
},
|
||||
origin_host=PROJECT_HOST,
|
||||
)
|
||||
except Exception as exc:
|
||||
make_log("Events", f"Failed to record content_indexed event: {exc}", level="warning")
|
||||
|
||||
await session.commit()
|
||||
return platform_found, seqno
|
||||
@@ -308,6 +324,21 @@ async def indexer_loop(memory, platform_found: bool, seqno: int) -> [bool, int]:
|
||||
updated=datetime.now()
|
||||
)
|
||||
session.add(onchain_stored_content)
|
||||
try:
|
||||
await record_event(
|
||||
session,
|
||||
'content_indexed',
|
||||
{
|
||||
'onchain_index': item_index,
|
||||
'content_hash': item_content_hash_str,
|
||||
'encrypted_cid': item_content_cid_str,
|
||||
'item_address': item_address.to_string(1, 1, 1),
|
||||
'owner_address': item_owner_address.to_string(1, 1, 1) if item_owner_address else None,
|
||||
},
|
||||
origin_host=PROJECT_HOST,
|
||||
)
|
||||
except Exception as exc:
|
||||
make_log("Events", f"Failed to record content_indexed event: {exc}", level="warning")
|
||||
await session.commit()
|
||||
make_log("Indexer", f"Item indexed: {item_content_hash_str}", level="info")
|
||||
last_known_index += 1
|
||||
|
||||
@@ -18,9 +18,12 @@ from app.core.models.wallet_connection import WalletConnection
|
||||
from app.core._keyboards import get_inline_keyboard
|
||||
from app.core.models._telegram import Wrapped_CBotChat
|
||||
from app.core.storage import db_session
|
||||
from app.core._config import CLIENT_TELEGRAM_API_KEY
|
||||
from app.core._config import CLIENT_TELEGRAM_API_KEY, CLIENT_TELEGRAM_BOT_USERNAME, PROJECT_HOST
|
||||
from app.core.models.user import User
|
||||
from app.core.models import StarsInvoice
|
||||
from app.core.events.service import record_event
|
||||
from app.core._secrets import hot_pubkey
|
||||
from base58 import b58encode
|
||||
import os
|
||||
import traceback
|
||||
|
||||
@@ -53,15 +56,42 @@ async def license_index_loop(memory, platform_found: bool, seqno: int) -> [bool,
|
||||
|
||||
if star_payment.amount == existing_invoice.amount:
|
||||
if not existing_invoice.paid:
|
||||
user = (await session.execute(select(User).where(User.id == existing_invoice.user_id))).scalars().first()
|
||||
existing_invoice.paid = True
|
||||
existing_invoice.paid_at = datetime.utcnow()
|
||||
existing_invoice.telegram_id = getattr(user, 'telegram_id', None)
|
||||
existing_invoice.payment_tx_id = getattr(star_payment, 'id', None)
|
||||
existing_invoice.payment_node_id = b58encode(hot_pubkey).decode()
|
||||
existing_invoice.payment_node_public_host = PROJECT_HOST
|
||||
existing_invoice.bot_username = CLIENT_TELEGRAM_BOT_USERNAME
|
||||
existing_invoice.is_remote = False
|
||||
await record_event(
|
||||
session,
|
||||
'stars_payment',
|
||||
{
|
||||
'invoice_id': existing_invoice.external_id,
|
||||
'content_hash': existing_invoice.content_hash,
|
||||
'amount': existing_invoice.amount,
|
||||
'user_id': existing_invoice.user_id,
|
||||
'telegram_id': existing_invoice.telegram_id,
|
||||
'bot_username': CLIENT_TELEGRAM_BOT_USERNAME,
|
||||
'type': existing_invoice.type,
|
||||
'payment_node': {
|
||||
'public_key': b58encode(hot_pubkey).decode(),
|
||||
'public_host': PROJECT_HOST,
|
||||
},
|
||||
'paid_at': existing_invoice.paid_at.isoformat() + 'Z' if existing_invoice.paid_at else None,
|
||||
'payment_tx_id': existing_invoice.payment_tx_id,
|
||||
},
|
||||
origin_host=PROJECT_HOST,
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
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
|
||||
)
|
||||
if user and 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
|
||||
)
|
||||
except BaseException as e:
|
||||
make_log("StarsProcessing", f"Local error: {e}" + '\n' + traceback.format_exc(), level="error")
|
||||
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
from .service import (
|
||||
record_event,
|
||||
store_remote_events,
|
||||
verify_event_signature,
|
||||
next_local_seq,
|
||||
upsert_cursor,
|
||||
prune_events,
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
'record_event',
|
||||
'store_remote_events',
|
||||
'verify_event_signature',
|
||||
'next_local_seq',
|
||||
'upsert_cursor',
|
||||
'prune_events',
|
||||
]
|
||||
@@ -0,0 +1,185 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, Dict, Iterable, List, Optional
|
||||
from uuid import uuid4
|
||||
|
||||
from base58 import b58decode, b58encode
|
||||
import nacl.signing
|
||||
from sqlalchemy import select, delete
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.core.logger import make_log
|
||||
from app.core._secrets import hot_pubkey, hot_seed
|
||||
from app.core.models import NodeEvent, NodeEventCursor
|
||||
|
||||
|
||||
LOCAL_PUBLIC_KEY = b58encode(hot_pubkey).decode()
|
||||
|
||||
|
||||
def _normalize_dt(value: Optional[datetime]) -> datetime:
|
||||
if value is None:
|
||||
return datetime.utcnow()
|
||||
if value.tzinfo is not None:
|
||||
return value.astimezone(timezone.utc).replace(tzinfo=None)
|
||||
return value
|
||||
|
||||
|
||||
def _parse_iso_dt(iso_value: Optional[str]) -> datetime:
|
||||
if not iso_value:
|
||||
return datetime.utcnow()
|
||||
try:
|
||||
parsed = datetime.fromisoformat(iso_value.replace('Z', '+00:00'))
|
||||
except Exception:
|
||||
return datetime.utcnow()
|
||||
return _normalize_dt(parsed)
|
||||
|
||||
|
||||
def _canonical_blob(data: Dict[str, Any]) -> bytes:
|
||||
return json.dumps(data, sort_keys=True, separators=(",", ":")).encode()
|
||||
|
||||
|
||||
def _sign_event(blob: Dict[str, Any]) -> str:
|
||||
signing_key = nacl.signing.SigningKey(hot_seed)
|
||||
signature = signing_key.sign(_canonical_blob(blob)).signature
|
||||
return b58encode(signature).decode()
|
||||
|
||||
|
||||
def verify_event_signature(event: Dict[str, Any]) -> bool:
|
||||
try:
|
||||
origin_key = event["origin_public_key"]
|
||||
signature = event["signature"]
|
||||
payload = {
|
||||
"origin_public_key": origin_key,
|
||||
"origin_host": event.get("origin_host"),
|
||||
"seq": event["seq"],
|
||||
"uid": event["uid"],
|
||||
"event_type": event["event_type"],
|
||||
"payload": event.get("payload") or {},
|
||||
"created_at": event.get("created_at"),
|
||||
}
|
||||
verify_key = nacl.signing.VerifyKey(b58decode(origin_key))
|
||||
verify_key.verify(_canonical_blob(payload), b58decode(signature))
|
||||
return True
|
||||
except Exception as exc:
|
||||
make_log("Events", f"Signature validation failed: {exc}", level="warning")
|
||||
return False
|
||||
|
||||
|
||||
async def next_local_seq(session: AsyncSession) -> int:
|
||||
result = await session.execute(
|
||||
select(NodeEvent.seq)
|
||||
.where(NodeEvent.origin_public_key == LOCAL_PUBLIC_KEY)
|
||||
.order_by(NodeEvent.seq.desc())
|
||||
.limit(1)
|
||||
)
|
||||
row = result.scalar_one_or_none()
|
||||
return int(row or 0) + 1
|
||||
|
||||
|
||||
async def record_event(
|
||||
session: AsyncSession,
|
||||
event_type: str,
|
||||
payload: Dict[str, Any],
|
||||
origin_host: Optional[str] = None,
|
||||
created_at: Optional[datetime] = None,
|
||||
) -> NodeEvent:
|
||||
seq = await next_local_seq(session)
|
||||
created_dt = _normalize_dt(created_at)
|
||||
event_body = {
|
||||
"origin_public_key": LOCAL_PUBLIC_KEY,
|
||||
"origin_host": origin_host,
|
||||
"seq": seq,
|
||||
"uid": uuid4().hex,
|
||||
"event_type": event_type,
|
||||
"payload": payload,
|
||||
"created_at": created_dt.replace(tzinfo=timezone.utc).isoformat().replace('+00:00', 'Z'),
|
||||
}
|
||||
signature = _sign_event(event_body)
|
||||
node_event = NodeEvent(
|
||||
origin_public_key=LOCAL_PUBLIC_KEY,
|
||||
origin_host=origin_host,
|
||||
seq=seq,
|
||||
uid=event_body["uid"],
|
||||
event_type=event_type,
|
||||
payload=payload,
|
||||
signature=signature,
|
||||
created_at=created_dt,
|
||||
status='local',
|
||||
)
|
||||
session.add(node_event)
|
||||
await session.flush()
|
||||
make_log("Events", f"Recorded local event {event_type} seq={seq}")
|
||||
return node_event
|
||||
|
||||
|
||||
async def upsert_cursor(session: AsyncSession, source_public_key: str, seq: int, host: Optional[str]):
|
||||
existing = (await session.execute(
|
||||
select(NodeEventCursor).where(NodeEventCursor.source_public_key == source_public_key)
|
||||
)).scalar_one_or_none()
|
||||
if existing:
|
||||
if seq > existing.last_seq:
|
||||
existing.last_seq = seq
|
||||
if host:
|
||||
existing.source_public_host = host
|
||||
else:
|
||||
cursor = NodeEventCursor(
|
||||
source_public_key=source_public_key,
|
||||
last_seq=seq,
|
||||
source_public_host=host,
|
||||
)
|
||||
session.add(cursor)
|
||||
await session.flush()
|
||||
|
||||
|
||||
async def store_remote_events(
|
||||
session: AsyncSession,
|
||||
events: Iterable[Dict[str, Any]],
|
||||
allowed_public_keys: Optional[set[str]] = None,
|
||||
) -> List[NodeEvent]:
|
||||
stored: List[NodeEvent] = []
|
||||
for event in events:
|
||||
if not verify_event_signature(event):
|
||||
continue
|
||||
origin_pk = event["origin_public_key"]
|
||||
if allowed_public_keys is not None and origin_pk not in allowed_public_keys:
|
||||
make_log("Events", f"Ignored event from untrusted node {origin_pk}", level="warning")
|
||||
continue
|
||||
seq = int(event["seq"])
|
||||
exists = (await session.execute(
|
||||
select(NodeEvent).where(
|
||||
NodeEvent.origin_public_key == origin_pk,
|
||||
NodeEvent.seq == seq,
|
||||
)
|
||||
)).scalar_one_or_none()
|
||||
if exists:
|
||||
continue
|
||||
created_dt = _parse_iso_dt(event.get("created_at"))
|
||||
received_dt = datetime.utcnow()
|
||||
node_event = NodeEvent(
|
||||
origin_public_key=origin_pk,
|
||||
origin_host=event.get("origin_host"),
|
||||
seq=seq,
|
||||
uid=event["uid"],
|
||||
event_type=event["event_type"],
|
||||
payload=event.get("payload") or {},
|
||||
signature=event["signature"],
|
||||
created_at=created_dt,
|
||||
status='recorded',
|
||||
received_at=received_dt,
|
||||
)
|
||||
session.add(node_event)
|
||||
stored.append(node_event)
|
||||
await upsert_cursor(session, origin_pk, seq, event.get("origin_host"))
|
||||
make_log("Events", f"Ingested remote event {event['event_type']} from {origin_pk} seq={seq}", level="debug")
|
||||
if stored:
|
||||
await session.flush()
|
||||
return stored
|
||||
|
||||
|
||||
async def prune_events(session: AsyncSession, max_age_days: int = 90):
|
||||
cutoff = datetime.utcnow() - timedelta(days=max_age_days)
|
||||
await session.execute(
|
||||
delete(NodeEvent).where(NodeEvent.created_at < cutoff)
|
||||
)
|
||||
@@ -11,6 +11,7 @@ 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.events import NodeEvent, NodeEventCursor
|
||||
from app.core.models.promo import PromoAction
|
||||
from app.core.models.tasks import BlockchainTask
|
||||
from app.core.models.content_v3 import (
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from sqlalchemy import (
|
||||
Column,
|
||||
Integer,
|
||||
BigInteger,
|
||||
String,
|
||||
DateTime,
|
||||
JSON,
|
||||
UniqueConstraint,
|
||||
)
|
||||
|
||||
from .base import AlchemyBase
|
||||
|
||||
|
||||
class NodeEvent(AlchemyBase):
|
||||
__tablename__ = 'node_events'
|
||||
__table_args__ = (
|
||||
UniqueConstraint('origin_public_key', 'seq', name='uq_node_events_origin_seq'),
|
||||
UniqueConstraint('uid', name='uq_node_events_uid'),
|
||||
)
|
||||
|
||||
id = Column(Integer, autoincrement=True, primary_key=True)
|
||||
origin_public_key = Column(String(128), nullable=False)
|
||||
origin_host = Column(String(256), nullable=True)
|
||||
seq = Column(BigInteger, nullable=False)
|
||||
uid = Column(String(64), nullable=False)
|
||||
event_type = Column(String(64), nullable=False)
|
||||
payload = Column(JSON, nullable=False, default=dict)
|
||||
signature = Column(String(512), nullable=False)
|
||||
created_at = Column(DateTime, nullable=False, default=datetime.utcnow)
|
||||
received_at = Column(DateTime, nullable=False, default=datetime.utcnow)
|
||||
applied_at = Column(DateTime, nullable=True)
|
||||
status = Column(String(32), nullable=False, default='recorded')
|
||||
|
||||
|
||||
class NodeEventCursor(AlchemyBase):
|
||||
__tablename__ = 'node_event_cursors'
|
||||
__table_args__ = (
|
||||
UniqueConstraint('source_public_key', name='uq_event_cursor_source'),
|
||||
)
|
||||
|
||||
id = Column(Integer, autoincrement=True, primary_key=True)
|
||||
source_public_key = Column(String(128), nullable=False)
|
||||
last_seq = Column(BigInteger, nullable=False, default=0)
|
||||
updated_at = Column(DateTime, nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
|
||||
source_public_host = Column(String(256), nullable=True)
|
||||
@@ -49,8 +49,15 @@ class StarsInvoice(AlchemyBase):
|
||||
|
||||
user_id = Column(Integer, ForeignKey('users.id'), nullable=True)
|
||||
content_hash = Column(String(256), nullable=True)
|
||||
telegram_id = Column(Integer, nullable=True)
|
||||
|
||||
invoice_url = Column(String(256), nullable=True)
|
||||
paid = Column(Boolean, nullable=False, default=False)
|
||||
paid_at = Column(DateTime, nullable=True)
|
||||
payment_tx_id = Column(String(256), nullable=True)
|
||||
payment_node_id = Column(String(128), nullable=True)
|
||||
payment_node_public_host = Column(String(256), nullable=True)
|
||||
bot_username = Column(String(128), nullable=True)
|
||||
is_remote = Column(Boolean, nullable=False, default=False)
|
||||
|
||||
created = Column(DateTime, nullable=False, default=datetime.utcnow)
|
||||
@@ -1,5 +1,5 @@
|
||||
from datetime import datetime
|
||||
from sqlalchemy import Column, Integer, String, BigInteger, DateTime, JSON
|
||||
from sqlalchemy import Column, Integer, String, BigInteger, DateTime, JSON, Boolean
|
||||
from sqlalchemy.orm import relationship
|
||||
|
||||
from app.core.auth_v1 import AuthenticationMixin as AuthenticationMixin_V1
|
||||
@@ -23,6 +23,7 @@ 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=dict)
|
||||
is_admin = Column(Boolean, nullable=False, default=False)
|
||||
|
||||
last_use = Column(DateTime, nullable=False, default=datetime.utcnow)
|
||||
updated = Column(DateTime, nullable=False, default=datetime.utcnow)
|
||||
|
||||
@@ -62,9 +62,15 @@ def verify_request(request, memory) -> Tuple[bool, str, str]:
|
||||
import nacl.signing
|
||||
vk = nacl.signing.VerifyKey(b58decode(node_id))
|
||||
sig = b58decode(sig_b58)
|
||||
msg = canonical_string(request.method, request.path, request.body or b"", ts, nonce, node_id)
|
||||
path = request.path
|
||||
query_string = getattr(request, 'query_string', None)
|
||||
if query_string:
|
||||
if not isinstance(query_string, str):
|
||||
query_string = query_string.decode() if isinstance(query_string, bytes) else str(query_string)
|
||||
if query_string:
|
||||
path = f"{path}?{query_string}"
|
||||
msg = canonical_string(request.method, path, request.body or b"", ts, nonce, node_id)
|
||||
vk.verify(msg, sig)
|
||||
return True, node_id, ""
|
||||
except Exception as e:
|
||||
return False, "", f"BAD_SIGNATURE: {e}"
|
||||
|
||||
Reference in new issue
Block a user