moebiusT7/book-ocr-studio
0
1import fcntl, os, signal, subprocess, sys, time, traceback, uuid2from pathlib import Path3from core import ROOT, read, save, prepare, review_page, status, export4from gpu_server import LocalGemma5from owned_process import OwnedPopen6 7from scheduler import gpu_plan, learn8from gpu_fallback import execute, is_oom9 10def stop(proc):11 if proc:12 try:os.killpg(proc.pid,signal.SIGTERM)13 except ProcessLookupError:pass14 try: proc.wait(timeout=10)15 except subprocess.TimeoutExpired:16 os.killpg(proc.pid,signal.SIGKILL);proc.wait()17 # Includes children left behind after the process leader exits.18 try:os.killpg(proc.pid,signal.SIGKILL)19 except ProcessLookupError:pass20 21def ocr_failure(proc):22 path=getattr(proc,'ocr_log_path',None)23 detail=''24 if path and path.exists():25 with path.open('rb') as f:26 f.seek(max(0,path.stat().st_size-16000));detail=f.read().decode(errors='replace')27 if is_oom(detail):raise RuntimeError('OCR ran out of VRAM (CUDA out of memory)')28 raise RuntimeError('The OCR process failed. Check the OCR attempt log.')29 30def review_ready(job, cfg, endpoint, proc=None, gpu=None):31 errors=[]32 for idx in cfg['selected']:33 folder=job/f'page-{idx+1:05d}'34 while not (folder/'marker.json').exists():35 if (job/'cancel').exists(): return errors36 if proc is None or proc.poll() is not None:37 if proc is not None and proc.returncode:ocr_failure(proc)38 raise RuntimeError(f'OCR did not finish for page {idx+1}. Check the log.')39 status(job,'parallel','Waiting for OCR pages',page=idx+1)40 time.sleep(.2)41 if (job/'cancel').exists(): return errors42 if (folder/'review.json').exists(): continue43 status(job,'parallel' if proc and proc.poll() is None else 'gemma',f"{'External model' if cfg.get('review_connector') else 'Gemma'} is reviewing page {idx+1} (stop takes effect after the response)")44 try:45 started=time.time()46 review_page(folder,cfg['model'],endpoint)47 (folder/'review-error.json').unlink(missing_ok=True)48 save(folder/'gemma-timing.json',dict(start=started,end=time.time(),endpoint=endpoint,gpu=gpu))49 except Exception as exc:50 error=dict(error=str(exc),page=idx+1)51 save(folder/'review-error.json',error);errors.append(error)52 if is_oom(exc):raise53 export(job)54 return errors55 56def attempt(job,cfg,plan,ocr_python,ocr_script):57 proc=None;errors=[]58 save(job/'execution.json',plan)59 env=os.environ.copy();env['CUDA_VISIBLE_DEVICES']=str(plan['ocr_gpu']);env['TORCH_DEVICE']='cuda'60 log_path=job/('ocr-attempt-'+uuid.uuid4().hex+'.log')61 def launch_ocr():62 with log_path.open('ab') as log:63 child=OwnedPopen([ocr_python,ocr_script,str(job)],env=env,stdout=log,stderr=log,start_new_session=True)64 child.ocr_log_path=log_path65 return child66 def admit(gpu, minimum, stage):67 from gpu_admission import wait_for_memory68 def observe(event):69 save(job/'gpu-admission.json',dict(stage=stage,time=time.time(),**event))70 status(job,'waiting_gpu',f"Waiting for GPU {gpu}: {event['free_mib']}MiB free / {minimum}MiB required ({stage})")71 wait_for_memory(gpu,minimum,lambda:(job/'cancel').exists(),observe)72 try:73 # Fresh admission after planning/preparation, and again at the sequential handoff.74 admit(plan['ocr_gpu'],14500 if plan['mode']=='shared' else 10000,'OCR')75 status(job,'ocr',f"{plan['mode']} mode: running OCR and review")76 if cfg['gemma'] and not cfg.get('review_connector') and plan['mode'] in {'dual','shared'}:77 if plan['mode']=='dual':78 admit(plan['gemma_gpu'],10000 if cfg['model']=='gemma4:12b-it-qat' else 15300,'Gemma (dual)')79 with LocalGemma(plan['gemma_gpu'],job,model=cfg['model']) as endpoint:80 proc=launch_ocr()81 try:82 errors=review_ready(job,cfg,endpoint,proc,gpu=plan['gemma_gpu'])83 if not (job/'cancel').exists():proc.wait(timeout=30)84 finally:stop(proc)85 else:86 proc=launch_ocr()87 while proc.poll() is None:88 if (job/'cancel').exists():stop(proc);return errors89 time.sleep(.2)90 if (job/'cancel').exists():return errors91 if proc.returncode:ocr_failure(proc)92 if cfg['gemma'] and plan['mode']=='sequential':93 stop(proc) # OCR exits and releases allocations before loading Gemma.94 if not cfg.get('review_connector'):95 admit(plan['gemma_gpu'],10000 if cfg['model']=='gemma4:12b-it-qat' else 15300,'Gemma after OCR exit')96 if cfg.get('review_connector'):97 errors=review_ready(job,cfg,'compat:'+cfg['review_connector'])98 else:99 with LocalGemma(plan['gemma_gpu'],job,model=cfg['model']) as endpoint:100 errors=review_ready(job,cfg,endpoint,gpu=plan['gemma_gpu'])101 return errors102 finally:103 stop(proc)104 export(job)105 106def run(job):107 proc=None108 with (job/'run.lock').open('a') as lock:109 try: fcntl.flock(lock,fcntl.LOCK_EX|fcntl.LOCK_NB)110 except BlockingIOError: return111 try:112 with (ROOT/'pipeline.lock').open('a') as shared:113 status(job,'queued','Queued for processing')114 while True:115 try: fcntl.flock(shared,fcntl.LOCK_EX|fcntl.LOCK_NB);break116 except BlockingIOError:117 if (job/'cancel').exists():status(job,'cancelled','Stopped');return118 time.sleep(1)119 cfg=read(job/'job.json')120 from capture_options import apply_options121 cfg=apply_options(job,cfg,ROOT,read,save)122 external=bool(cfg.get('review_connector'))123 if external:124 from review_connector import model_for125 if model_for(cfg['review_connector'])!=cfg['model']:raise ValueError('Connector model mismatch')126 engine=cfg.get('ocr_engine','marker')127 if engine not in {'marker','yomitoku'}:raise RuntimeError('Unknown OCR engine')128 ocr_python=str(ROOT/'.venv-yomitoku/bin/python') if engine=='yomitoku' else sys.executable129 ocr_script=str(ROOT/('yomitoku_worker.py' if engine=='yomitoku' else 'marker_worker.py'))130 if not Path(ocr_python).exists():raise RuntimeError('YomiToku environment not found')131 from calibrate_gpu import calibrate132 status(job,'calibrating','Checking transfer bandwidth on available GPUs')133 save(job/'calibration.json',calibrate())134 history_kind=cfg['kind']+(':yomitoku' if engine=='yomitoku' else '')+(':'+cfg['model'] if cfg['model']!='gemma4:26b-a4b-it-qat' else '')135 plan=gpu_plan('sequential' if external else cfg.get('mode','auto'),cfg['gemma'] and not external,len(cfg['selected']),history_kind,model=cfg['model']);save(job/'execution.json',plan)136 status(job,'preparing','Saving source pages')137 prepare(job)138 if (job/'cancel').exists():status(job,'cancelled','Stopped');return139 def select(mode):140 # Owned processes have exited before fresh GPU inventory is read.141 return gpu_plan('sequential' if external else mode,cfg['gemma'] and not external,len(cfg['selected']),history_kind,model=cfg['model'])142 def record(event):143 path=job/'fallback.json'144 history=read(path) if path.exists() else []145 history.append(dict(time=time.time(),**event));save(path,history)146 status(job,'switching',f"VRAM fallback: {event['event']} → {event.get('mode')}")147 errors,plan=execute(plan,lambda p:attempt(job,cfg,p,ocr_python,ocr_script),select,record,lambda:(job/'cancel').exists())148 if (job/'cancel').exists():status(job,'cancelled','Stopped. You can resume from saved pages.');return149 export(job)150 from delivery import ensure_delivery151 ensure_delivery(job)152 if not external and engine=='marker' and cfg['model']=='gemma4:26b-a4b-it-qat':learn(job,plan)153 status(job,'partial' if errors else 'done','OCR complete. Some model reviews failed.' if errors else 'Processing complete. Suggestions have not been approved.',errors=errors)154 except Exception as exc:155 stop(proc);traceback.print_exc()156 try:export(job)157 except Exception:pass158 if (job/'cancel').exists():status(job,'cancelled','Stopped. Saved results are preserved.')159 else:status(job,'error',str(exc))160 finally:stop(proc)161if __name__=='__main__':162 os.umask(0o077)163 run(Path(sys.argv[1]).resolve())164 