CoolFace
Apppublic

moebiusT7/book-ocr-studio

sourceHugging Faceagpl-3.0updated 3d agoView on Hugging Face
0likes
worker.py164 linesDownload Raw Back to root
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