automatic handshake and connect
This commit is contained in:
1 parent
dbc460f0bb
commit
da446f5ab0
7 files changed
+283
-10
No files matched your search
@@ -3,6 +3,7 @@ import os
|
||||
from typing import List, Optional
|
||||
|
||||
import httpx
|
||||
from urllib.parse import urlparse
|
||||
import random
|
||||
import shutil
|
||||
from sqlalchemy import select
|
||||
@@ -11,7 +12,7 @@ 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.content_v3 import EncryptedContent, ContentDerivative
|
||||
from app.core.ipfs_client import pin_add, find_providers, swarm_connect
|
||||
from app.core.ipfs_client import pin_add, pin_ls, find_providers, swarm_connect, add_streamed_file
|
||||
|
||||
|
||||
INTERVAL_SEC = 60
|
||||
@@ -28,7 +29,8 @@ async def fetch_index(base_url: str, etag: Optional[str], since: Optional[str])
|
||||
url = f"{base_url.rstrip('/')}/api/v1/content.delta" if since else f"{base_url.rstrip('/')}/api/v1/content.index"
|
||||
if etag:
|
||||
headers['If-None-Match'] = etag
|
||||
async with httpx.AsyncClient(timeout=20) as client:
|
||||
# follow_redirects handles peers that force HTTPS and issue 301s
|
||||
async with httpx.AsyncClient(timeout=20, follow_redirects=True) as client:
|
||||
r = await client.get(url, headers=headers, params=params)
|
||||
if r.status_code != 200:
|
||||
if r.status_code == 304:
|
||||
@@ -157,6 +159,20 @@ async def main_fn(memory):
|
||||
async def _pin_one(cid: str):
|
||||
async with sem:
|
||||
try:
|
||||
node_ipfs_meta = (n.meta or {}).get('ipfs') or {}
|
||||
multiaddrs = node_ipfs_meta.get('multiaddrs') or []
|
||||
for addr in multiaddrs:
|
||||
try:
|
||||
await swarm_connect(addr)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
existing = await pin_ls(cid)
|
||||
if existing and existing.get('Keys'):
|
||||
make_log('index_scout_v3', f"pin {cid} already present", level='debug')
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
# Try to pre-connect to discovered providers
|
||||
try:
|
||||
provs = await find_providers(cid, max_results=5)
|
||||
@@ -168,15 +184,56 @@ async def main_fn(memory):
|
||||
pass
|
||||
except Exception:
|
||||
pass
|
||||
await pin_add(cid, recursive=True)
|
||||
try:
|
||||
await asyncio.wait_for(pin_add(cid, recursive=True), timeout=60)
|
||||
return
|
||||
except httpx.HTTPStatusError as http_err:
|
||||
body = (http_err.response.text or '').lower() if http_err.response else ''
|
||||
if 'already pinned' in body or 'pin already set' in body:
|
||||
make_log('index_scout_v3', f"pin {cid} already present", level='debug')
|
||||
return
|
||||
raise
|
||||
except Exception as e:
|
||||
make_log('index_scout_v3', f"pin {cid} failed: {e}", level='warning')
|
||||
# Attempt HTTP gateway fallback before logging failure
|
||||
fallback_sources = []
|
||||
node_host = n.meta.get('public_host') if isinstance(n.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_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}"
|
||||
make_log('index_scout_v3', f"fallback download start {cid} via {gateway_url}", level='debug')
|
||||
async with httpx.AsyncClient(timeout=None) as client:
|
||||
resp = await client.get(gateway_url)
|
||||
resp.raise_for_status()
|
||||
data = resp.content
|
||||
chunk_bytes = int(os.getenv('CRYPTO_CHUNK_BYTES', '1048576'))
|
||||
add_params = {
|
||||
'cid-version': 1,
|
||||
'raw-leaves': 'true',
|
||||
'chunker': f'size-{chunk_bytes}',
|
||||
'hash': 'sha2-256',
|
||||
'pin': 'true',
|
||||
}
|
||||
result = await add_streamed_file([data], filename=f'{cid}.bin', params=add_params)
|
||||
if str(result.get('Hash')) != str(cid):
|
||||
raise ValueError(f"gateway add returned mismatched CID {result.get('Hash')}")
|
||||
make_log('index_scout_v3', f"pin {cid} fetched via gateway {gateway_host}:{gateway_port}", level='info')
|
||||
return
|
||||
else:
|
||||
fallback_sources.append('gateway-host-missing')
|
||||
except Exception as fallback_err:
|
||||
fallback_sources.append(str(fallback_err))
|
||||
make_log('index_scout_v3', f"pin {cid} failed: {e}; fallback={'; '.join(fallback_sources) if fallback_sources else 'none'}", level='warning')
|
||||
|
||||
tasks = []
|
||||
for it in items:
|
||||
await upsert_content(it)
|
||||
cid = it.get('encrypted_cid')
|
||||
if cid:
|
||||
make_log('index_scout_v3', f"queue pin {cid}")
|
||||
tasks.append(asyncio.create_task(_pin_one(cid)))
|
||||
if tasks:
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
Reference in new issue
Block a user