Files
SHiNE-server/tools/test-publisher/v2/publisher/publisher.py
T
AidarKC de8700f8e0 back32hash 5591: готово к тестам
Вроде как всё готово к тестам.
2026-10-02 01:32:37 +03:00

196 lines
11 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
import argparse, asyncio, base64, hashlib, json, os, random, struct, sys, time, uuid
from pathlib import Path
import re
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from websockets.asyncio.client import connect
ZERO32 = bytes(32)
FRAME_CODE_V1=1; TEXT_TYPE=1; TEXT_POST=10; VERSION=1
def shine_protocol_namespace():
here = Path(__file__).resolve()
for parent in here.parents:
cfg = parent / 'shine-UI' / 'js' / 'services' / 'protocol-tags.js'
if cfg.exists():
m = re.search(r"SHINE_PROTOCOL_NAMESPACE\s*=\s*['\"]([^'\"]+)['\"]", cfg.read_text(encoding='utf-8'))
if m:
return m.group(1)
raise RuntimeError('Не найден shine-UI/js/services/protocol-tags.js с SHINE_PROTOCOL_NAMESPACE')
SHINE_PROTOCOL_NAMESPACE = shine_protocol_namespace()
def atomic_json(path, obj):
path=Path(path); path.parent.mkdir(parents=True, exist_ok=True)
tmp=path.with_suffix(path.suffix+'.tmp')
tmp.write_text(json.dumps(obj,ensure_ascii=False,indent=2),encoding='utf-8')
os.replace(tmp,path)
def load_json(path, default=None):
p=Path(path)
return json.loads(p.read_text(encoding='utf-8')) if p.exists() else default
def avro_long(n):
n=(n<<1) ^ (n>>63); out=bytearray()
while n & ~0x7f: out.append((n&0x7f)|0x80); n >>= 7
out.append(n); return bytes(out)
def avro_tags(tags):
out=bytearray()
if tags:
out += avro_long(len(tags))
for k,v in tags:
kb=k.encode(); vb=v.encode(); out+=avro_long(len(kb))+kb+avro_long(len(vb))+vb
out += avro_long(0); return bytes(out)
def deep_hash(x):
H=lambda b: hashlib.sha384(b).digest()
if isinstance(x,(bytes,bytearray)):
b=bytes(x); return H(H(f'blob{len(b)}'.encode())+H(b))
acc=H(f'list{len(x)}'.encode())
for child in x: acc=H(acc+deep_hash(child))
return acc
def signing_message(owner,tags,data):
raw=avro_tags(tags)
return deep_hash([b'dataitem',b'1',b'2',owner,b'',b'',raw,data])
def data_item(priv64,tags,data):
if len(priv64) not in (32,64): raise ValueError('Solana key JSON must contain 32 or 64 bytes')
seed=priv64[:32]; owner=priv64[32:64] if len(priv64)==64 else Ed25519PrivateKey.from_private_bytes(seed).public_key().public_bytes_raw()
sig=Ed25519PrivateKey.from_private_bytes(seed).sign(signing_message(owner,tags,data))
rawtags=avro_tags(tags)
return struct.pack('<H',2)+sig+owner+b'\0\0'+struct.pack('<QQ',len(tags),len(rawtags))+rawtags+data
def post_body(line_code, prev_line_num, prev_line_hash, this_line_num, text):
tb=text.encode('utf-8')
if len(tb)>65535: raise ValueError('post text too long')
return struct.pack('>ii32siH',line_code,prev_line_num,prev_line_hash,this_line_num,len(tb))+tb
def frame(prev_hash, back32_hash, block_num, body, ts=None):
ts=int(time.time()) if ts is None else int(ts); size=88+len(body)
return struct.pack('>H32s32siiqHHH',FRAME_CODE_V1,prev_hash,back32_hash,size,block_num,ts,TEXT_TYPE,TEXT_POST,VERSION)+body
def h32(b): return hashlib.sha256(b).digest()
def hx(b): return b.hex()
async def ws_call(url, op, payload, timeout=20):
req={'op':op,'requestId':str(uuid.uuid4()),'payload':payload}
async with connect(url, open_timeout=timeout, close_timeout=5) as ws:
await ws.send(json.dumps(req,separators=(',',':')))
end=time.monotonic()+timeout
while True:
left=end-time.monotonic()
if left<=0: raise TimeoutError(op)
msg=await asyncio.wait_for(ws.recv(),left)
obj=json.loads(msg)
if obj.get('requestId')==req['requestId']: return obj
async def head(cfg):
r=await ws_call(cfg['server_ws'],'ListBlockchainHeads',{})
if r.get('status')!=200: raise RuntimeError(f'ListBlockchainHeads: {r}')
for x in r.get('payload',{}).get('blockchains',[]):
if x.get('blockchainName')==cfg['blockchain_name']:
return int(x.get('lastBlockNumber',-1)), x.get('lastBlockHash','')
raise RuntimeError('blockchain not found on server')
async def block(cfg,n):
r=await ws_call(cfg['server_ws'],'GetBlockchainBlock',{'blockchainName':cfg['blockchain_name'],'blockNumber':n})
if r.get('status')!=200: raise RuntimeError(f'GetBlockchainBlock({n}): {r}')
p=r.get('payload',r)
return p.get('blockHash') or r.get('blockHash'), p.get('blockBytesB64') or r.get('blockBytesB64')
async def channel_tail(cfg, root_hash):
r=await ws_call(cfg['server_ws'],'GetChannelMessages',{'channel':{'ownerBlockchainName':cfg['blockchain_name'],'channelRootBlockNumber':cfg['channel_root_block_number'],'channelRootBlockHash':root_hash},'limit':1,'sort':'desc'})
if r.get('status')!=200: raise RuntimeError(f'GetChannelMessages: {r}')
items=r.get('payload',{}).get('messages',[])
if not items: return 0, ZERO32
m=items[0]; ref=m.get('messageRef') or {}
return int(m.get('lineStep') or 0), bytes.fromhex(ref.get('blockHash','00'*32))
def read_key(path):
a=load_json(path)
if not isinstance(a,list): raise ValueError('key file must be Solana JSON byte array')
return bytes(int(x)&255 for x in a)
def queue_items(cfg):
q=load_json(cfg['queue_file'])
if not isinstance(q,list) or not q: raise ValueError('queue must be non-empty JSON array')
for x in q:
if not isinstance(x,dict) or not str(x.get('text','')).strip(): raise ValueError('each queue item needs text')
return q
def delay(cfg): return random.randint(int(cfg['min_interval_seconds']),int(cfg['max_interval_seconds']))
async def init_state(cfg):
bn,bhash=await head(cfg)
root_hash,_=await block(cfg,int(cfg['channel_root_block_number']))
line_num,line_hash=await channel_tail(cfg,root_hash)
st={'version':1,'next_index':0,'last_block_number':bn,'last_block_hash':bhash,'last_line_number':line_num,'last_line_hash':hx(line_hash),'next_publish_at':int(time.time())+delay(cfg),'pending':None,'published_total':0}
atomic_json(cfg['state_file'],st); return st
async def reconcile(cfg,st):
p=st.get('pending')
bn,bhash=await head(cfg)
if p:
if bn==p['block_number'] and bhash.lower()==p['block_hash'].lower():
st['last_block_number']=bn; st['last_block_hash']=bhash; st['last_line_number']=p['this_line_number']; st['last_line_hash']=p['block_hash']; st['next_index']=p['next_index_after']; st['published_total']=int(st.get('published_total',0))+1; st['pending']=None; st['next_publish_at']=int(time.time())+delay(cfg); atomic_json(cfg['state_file'],st); print('Recovered accepted pending block',bn)
elif bn==p['previous_block_number'] and bhash.lower()==p['previous_block_hash'].lower():
print('Pending block was not accepted; it will be retried')
else: raise RuntimeError(f'chain changed while pending: server={bn}:{bhash} state={p}')
else:
if bn!=st['last_block_number'] or bhash.lower()!=str(st['last_block_hash']).lower(): raise RuntimeError('server chain head differs from publisher state; refusing to fork')
return st
async def publish_one(cfg,st,items):
if st.get('pending'):
p=st['pending']; raw=base64.b64decode(p['data_item_b64']); block_num=p['block_number']; prev=p['previous_block_hash']
else:
idx=int(st['next_index'])
if idx>=len(items):
if cfg.get('loop_queue',False): idx=0
else: print('Queue exhausted. Nothing to publish.'); return False
text=items[idx]['text'].strip(); block_num=int(st['last_block_number'])+1; prev_hash=bytes.fromhex(st['last_block_hash']); this_line=int(st['last_line_number'])+1
body=post_body(int(cfg['channel_root_block_number']),int(st['last_line_number']),bytes.fromhex(st['last_line_hash']),this_line,text)
back32_hash=ZERO32 if block_num<32 else bytes.fromhex((await block(cfg,block_num-32))[0]); fr=frame(prev_hash,back32_hash,block_num,body); block_hash=h32(fr); tags=[('App',SHINE_PROTOCOL_NAMESPACE),('c',cfg['channel_name'].strip().lower())]
raw=data_item(read_key(cfg['key_file']),tags,fr)
nextidx=idx+1
p={'queue_index':idx,'next_index_after':nextidx,'block_number':block_num,'block_hash':hx(block_hash),'previous_block_number':st['last_block_number'],'previous_block_hash':st['last_block_hash'],'this_line_number':this_line,'data_item_b64':base64.b64encode(raw).decode()}
st['pending']=p; atomic_json(cfg['state_file'],st); prev=st['last_block_hash']
r=await ws_call(cfg['server_ws'],'AddBlock',{'blockchainName':cfg['blockchain_name'],'blockNumber':block_num,'prevBlockHash':prev,'blockBytesB64':base64.b64encode(raw).decode()},timeout=30)
if r.get('status')!=200: raise RuntimeError(f'AddBlock failed: {r}')
server_hash=(r.get('payload') or {}).get('serverLastBlockHash') or (r.get('payload') or {}).get('serverLastGlobalHash')
if server_hash and server_hash.lower()!=p['block_hash'].lower(): raise RuntimeError('server accepted different hash')
st['last_block_number']=block_num; st['last_block_hash']=p['block_hash']; st['last_line_number']=p['this_line_number']; st['last_line_hash']=p['block_hash']; st['next_index']=p['next_index_after']; st['published_total']=int(st.get('published_total',0))+1; st['pending']=None; st['next_publish_at']=int(time.time())+delay(cfg); atomic_json(cfg['state_file'],st)
print(f"Published queue[{p['queue_index']}] as block #{block_num}; next at {time.strftime('%Y-%m-%d %H:%M:%S',time.localtime(st['next_publish_at']))}")
return True
async def main_async(args):
cfg=load_json(args.config); required=['server_ws','blockchain_name','channel_name','channel_root_block_number','key_file','queue_file','state_file','min_interval_seconds','max_interval_seconds']
for k in required:
if k not in cfg: raise ValueError('missing config: '+k)
base=Path(args.config).resolve().parent
for k in ('key_file','queue_file','state_file'):
p=Path(cfg[k]); cfg[k]=str(p if p.is_absolute() else base/p)
items=queue_items(cfg)
st=load_json(cfg['state_file'])
if st is None: st=await init_state(cfg); print('Initialized state from server/channel')
st=await reconcile(cfg,st)
if args.command=='status': print(json.dumps(st,ensure_ascii=False,indent=2)); return
if args.command=='publish-now': await publish_one(cfg,st,items); return
while True:
st=await reconcile(cfg,load_json(cfg['state_file']))
wait=max(0,int(st.get('next_publish_at',0))-int(time.time()))
if wait: print(f'Waiting {wait}s'); await asyncio.sleep(wait)
ok=await publish_one(cfg,st,items)
if not ok: await asyncio.sleep(3600)
def main():
ap=argparse.ArgumentParser(description='SHiNE test channel publisher')
ap.add_argument('command',choices=['run','status','publish-now'],nargs='?',default='run'); ap.add_argument('--config',default='config.json')
a=ap.parse_args()
try: asyncio.run(main_async(a))
except KeyboardInterrupt: pass
except Exception as e: print('ERROR:',e,file=sys.stderr); sys.exit(1)
if __name__=='__main__': main()