File size: 15,656 Bytes
7d004b5
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
99bfb2e
 
7d004b5
 
 
 
 
 
99bfb2e
7d004b5
 
99bfb2e
7d004b5
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
"""One model per process. Save each complete cell and its document predictions.

Scores are strict character-boundary + type micro F1. No tuning on test labels.
"""
import argparse, contextlib, gc, hashlib, json, os, re, time, traceback
from collections import defaultdict
from pathlib import Path
import numpy as np
import torch
from huggingface_hub import snapshot_download

ROOT=Path(__file__).resolve().parents[1]

def dump(path,obj):
    path.parent.mkdir(parents=True,exist_ok=True)
    tmp=path.with_suffix(path.suffix+'.tmp');tmp.write_text(json.dumps(obj,indent=2,ensure_ascii=False));tmp.replace(path)

def chunks(text,size=64,overlap=16):
    words=list(re.finditer(r'\S+',text))
    if not words:return []
    out=[]
    for first in range(0,len(words),size-overlap):
        last=min(first+size,len(words));s=words[first].start();e=words[last-1].end()
        out.append((s,text[s:e]))
        if last==len(words):break
    return out

def counts(gold,pred):
    g={tuple(x) for x in gold};p={tuple(x) for x in pred}
    return [len(g&p),len(p-g),len(g-p)]

def metric(c):
    tp,fp,fn=map(int,c)
    return {'tp':tp,'fp':fp,'fn':fn,'precision':tp/(tp+fp) if tp+fp else 0.,
            'recall':tp/(tp+fn) if tp+fn else 0.,'f1':2*tp/(2*tp+fp+fn) if 2*tp+fp+fn else 0.}

def summarize(rows):
    arr=np.array([r['counts'] for r in rows],dtype=np.int64).reshape(-1,3)
    out=metric(arr.sum(axis=0));out['documents']=len(rows)
    if len(rows)>1:
        rng=np.random.default_rng(20260928)
        boot=arr[rng.integers(0,len(rows),(1000,len(rows)))].sum(axis=1)
        den=2*boot[:,0]+boot[:,1]+boot[:,2]
        vals=np.divide(2*boot[:,0],den,out=np.zeros(len(boot)),where=den!=0)
        out['f1_95_ci_document_bootstrap']=[float(x) for x in np.quantile(vals,[.025,.975])]
    return out

class Adapter:
    def __init__(self,entry):
        self.entry=entry;self.kind=entry['kind'];self.events=defaultdict(int)
        os.environ['TOKENIZERS_PARALLELISM']='false'
        torch.set_num_threads(4)
        path=snapshot_download(entry['id'],revision=entry['revision'],ignore_patterns=['*.bin','*.onnx','*.h5','*.msgpack','*.ot'])
        # Some T5 repositories contain only a PyTorch .bin checkpoint.
        if not list(Path(path).glob('*.safetensors')):
            path=snapshot_download(entry['id'],revision=entry['revision'])
        self.path=path
        if self.kind=='gliner2':
            from gliner2 import GLiNER2
            self.model=GLiNER2.from_pretrained(path).to('cuda').eval();self.batch=8
            self.max_len=512 if 'specialised' in entry['id'] else 3072
            core=getattr(self.model,'model',self.model)
            proc=getattr(core,'processor',getattr(self.model,'processor',None))
            self.runtime={'processor_max_length':getattr(proc,'max_length',None),'evaluation_max_len_argument':self.max_len}
        elif self.kind=='gliner':
            from gliner import GLiNER
            self.model=GLiNER.from_pretrained(path).to('cuda').eval();self.batch=8
            self.runtime={'description_interface':'label: definition strings mapped back to canonical labels'}
        elif self.kind=='gner':
            from transformers import AutoTokenizer,AutoModelForSeq2SeqLM
            self.tokenizer=AutoTokenizer.from_pretrained(path)
            self.model=AutoModelForSeq2SeqLM.from_pretrained(path,dtype=torch.bfloat16).to('cuda').eval();self.batch=8
            self.runtime={'precision':'bfloat16','generation':'greedy; max_new_tokens=512; official BIO prompt plus concise definitions'}
        elif self.kind=='nuextract':
            from transformers import AutoTokenizer,Qwen2_5_VLForConditionalGeneration
            self.tokenizer=AutoTokenizer.from_pretrained(path,padding_side='left')
            self.model=Qwen2_5_VLForConditionalGeneration.from_pretrained(path,dtype=torch.bfloat16).to('cuda').eval();self.batch=8
            self.runtime={'precision':'bfloat16','generation':'greedy; max_new_tokens=512; native template with verbatim-string arrays'}
        else:raise ValueError(self.kind)
        self.runtime.update({'device':torch.cuda.get_device_name(),'parameter_count':sum(p.numel() for p in self.model.parameters())})

    def predict(self,texts,schema):
        with torch.inference_mode():
            if self.kind=='gliner2':
                raw=self.model.batch_extract_entities(texts,schema,batch_size=self.batch,threshold=.5,include_confidence=True,include_spans=True,max_len=self.max_len)
                return [[[x['start'],x['end'],label] for label,values in row.get('entities',{}).items() for x in values] for row in raw]
            if self.kind=='gliner':
                mapping={f'{k}: {v}':k for k,v in schema.items()}
                raw=self.model.batch_predict_entities(texts,list(mapping),threshold=.5,flat_ner=True,batch_size=self.batch)
                return [[[x['start'],x['end'],mapping.get(x['label'],x['label'])] for x in row] for row in raw]
            if self.kind=='gner':
                instruction=("Please analyze the sentence provided, identifying the type of entity for each word on a token-by-token basis.\n"
                    "Output format is: word_1(label_1), word_2(label_2), ...\nWe'll use the BIO-format to label the entities, where:\n"
                    "1. B- (Begin) indicates the start of a named entity.\n2. I- (Inside) is used for words within a named entity but are not the first word.\n"
                    "3. O (Outside) denotes words that are not part of a named entity.\n")
                header=instruction+'\nUse the specific entity tags: '+', '.join(schema)+' and O.\nEntity definitions: '+json.dumps(schema)+'.\nSentence: '
                encoded=self.tokenizer([header+t for t in texts],return_tensors='pt',padding=True,truncation=False).to('cuda')
            else:
                template=json.dumps({k:['verbatim-string'] for k in schema},ensure_ascii=False)
                prompts=[]
                for text in texts:
                    # Definitions are metadata; only the delimited source is eligible for extraction.
                    content='Entity definitions (instructions, not source text): '+json.dumps(schema)+'\nExtract only from the following source text:\n'+text
                    messages=[{'role':'user','content':[{'type':'text','text':content}]}]
                    prompts.append(self.tokenizer.apply_chat_template(messages,template=template,tokenize=False,add_generation_prompt=True))
                encoded=self.tokenizer(prompts,return_tensors='pt',padding=True,truncation=False).to('cuda')
            outputs=self.model.generate(**encoded,max_new_tokens=512,do_sample=False)
            if self.kind=='nuextract':outputs=outputs[:,encoded['input_ids'].shape[1]:]
            result=[]
            for text,tokens in zip(texts,outputs):
                if len(tokens)>=512 and tokens[-1].item()!=self.tokenizer.eos_token_id:self.events['generation_limit_reached']+=1
                response=self.tokenizer.decode(tokens,skip_special_tokens=True)
                if self.events['raw_output_samples_saved']<12:
                    folder=ROOT/'diagnostics';folder.mkdir(exist_ok=True)
                    with (folder/(self.entry['id'].replace('/','--')+'.jsonl')).open('a') as stream:
                        stream.write(json.dumps({'text':text,'schema':schema,'output':response},ensure_ascii=False)+'\n')
                    self.events['raw_output_samples_saved']+=1
                if self.kind=='gner':result.append(self.parse_gner(text,response,schema))
                else:result.append(self.parse_json(text,response,schema))
            return result

    def parse_gner(self,text,response,schema):
        # Align BIO output to consecutive source tokens; retain unaligned entity predictions as FPs.
        # The author's model emits whitespace-separated tagged words even though
        # the prompt illustrates commas. Parse tags independently of separators.
        pairs=[];previous_end=0
        for match in re.finditer(r'\((B-[^()]+|I-[^()]+|O)\)',response):
            word=response[previous_end:match.start()].strip()
            if word.startswith(', '):word=word[2:].strip()
            pairs.append((word,match.group(1)));previous_end=match.end()
        source=list(re.finditer(r'\S+',text));cursor=0;ents=[];active=None;unknown=0
        aliases={k.casefold():k for k in schema}
        for word,tag in pairs:
            word=word.strip();prefix,_,label=tag.partition('-');label=aliases.get(label.casefold(),label)
            found=None
            for i in range(cursor,min(cursor+12,len(source))):
                if source[i].group()==word:found=i;break
            if found is not None:
                s,e=source[found].span();cursor=found+1
            else:s,e=-1,-1;self.events['unaligned_output_tokens']+=1
            if active and (prefix!='I' or label!=active[2] or s<0 or any(not c.isspace() for c in text[active[1]:s])):
                ents.append(active);active=None
            if prefix in ('B','I'):
                if s<0:
                    unknown+=1;ents.append([-100000-unknown,-100000-unknown,label]);continue
                if active:active[1]=e
                else:active=[s,e,label]
        if active:ents.append(active)
        if not pairs:self.events['unparseable_outputs']+=1
        return ents

    def parse_json(self,text,response,schema):
        response=response.strip();start=response.find('{');end=response.rfind('}')
        try:obj=json.loads(response[start:end+1])
        except Exception:
            self.events['invalid_json_outputs']+=1
            return [[-999999,-999999,'__invalid_output__']]
        result=[];used=set();bad=0
        for label,values in obj.items():
            if isinstance(values,str):values=[values]
            if not isinstance(values,list):self.events['invalid_value_types']+=1;continue
            for value in values:
                if not isinstance(value,str) or not value:continue
                found=next((m for m in re.finditer(re.escape(value),text) if (m.start(),m.end(),label) not in used),None)
                if found:
                    used.add((found.start(),found.end(),label));result.append([found.start(),found.end(),label])
                else:
                    bad+=1;result.append([-100000-bad,-100000-bad,label]);self.events['unaligned_strings']+=1
        return result

def evaluate(adapter,name,meta,track,docs,exclude,deadline):
    root=ROOT/'results'/adapter.entry['id'].replace('/','--');path=root/f'{name}.{track}.json'
    if path.exists() and json.loads(path.read_text()).get('complete'):return
    decisions=meta['labels']
    if track=='benchmark_schema':
        schema={k:v['definition'] for k,v in decisions.items()};mapping={k:k for k in schema}
    else:
        mapping={k:v['training_label'] for k,v in decisions.items() if v['status']=='familiar_equivalent'}
        schema={v['training_label']:v['training_definition'] for v in decisions.values() if v['status']=='familiar_equivalent'}
    if not schema:
        dump(path,{'complete':True,'not_applicable':True,'dataset':name,'track':track});return
    rows=[];start=time.time();model_seconds=0;total_chunks=0
    for base in range(0,len(docs),8):
        if time.time()>deadline:raise TimeoutError('Evaluation deadline reached')
        batch=docs[base:base+8];parts=[];owner=[]
        for i,d in enumerate(batch):
            for offset,txt in chunks(d['text']):parts.append(txt);owner.append((i,offset))
        predictions=[[] for _ in batch]
        for at in range(0,len(parts),8):
            tick=time.time()
            request=parts[at:at+8]
            try:
                outputs=adapter.predict(request,schema) if adapter.batch>1 else [adapter.predict([txt],schema)[0] for txt in request]
            except torch.cuda.OutOfMemoryError:
                torch.cuda.empty_cache();adapter.batch=1
                outputs=[adapter.predict([txt],schema)[0] for txt in request]
            model_seconds+=time.time()-tick
            if len(outputs)!=len(request):
                outputs=[adapter.predict([txt],schema)[0] for txt in request]
            assert len(outputs)==len(request)
            for j,values in enumerate(outputs):
                i,offset=owner[at+j]
                for s,e,label in values:
                    if s>=0:
                        if not 0<=s<e<=len(parts[at+j]):
                            adapter.events['invalid_offsets']+=1;predictions[i].append([-900000-at-j,-900000-at-j,label]);continue
                        s+=offset;e+=offset
                    predictions[i].append([s,e,label])
        total_chunks+=len(parts)
        for d,preds in zip(batch,predictions):
            gold=[[s,e,mapping[label]] for s,e,label in d['entities'] if label in mapping]
            preds=[list(x) for x in sorted(set(tuple(x) for x in preds))]
            row={'id':d['id'],'counts':counts(gold,preds),'gold':gold,'predictions':preds,'overlap_flag':d['id'] in exclude}
            if track=='benchmark_schema':
                novel={k for k,v in decisions.items() if v['status']=='new_specific_type_related'}
                row['new_specific_type_counts']=counts([x for x in gold if x[2] in novel],[x for x in preds if x[2] in novel])
            rows.append(row)
        dump(root/f'{name}.{track}.progress.json',{'documents_completed':len(rows),'documents_total':len(docs),'time':time.time()})
    clean=[r for r in rows if not r['overlap_flag']]
    result={'complete':True,'model':adapter.entry['id'],'revision':adapter.entry['revision'],'dataset':name,'track':track,
        'protocol':'fixed64-v1','test_documents':len(docs),'unfiltered':summarize(rows),'overlap_screened':summarize(clean),
        'latency_seconds':time.time()-start,'inference_seconds':model_seconds,'text_chunks':total_chunks,
        'runtime':adapter.runtime,'adapter_events_cumulative':dict(adapter.events),'schema':schema,'documents':rows}
    if track=='benchmark_schema':
        result['new_specific_types_related_to_training']=summarize([{'counts':r['new_specific_type_counts']} for r in clean])
    dump(path,result);print('RESULT',adapter.entry['id'],name,track,round(result['overlap_screened']['f1']*100,2),flush=True)

def main():
    p=argparse.ArgumentParser();p.add_argument('--model',required=True);p.add_argument('--deadline',type=float,required=True)
    p.add_argument('--datasets',nargs='*');p.add_argument('--tracks',nargs='*');args=p.parse_args()
    entries=json.loads((ROOT/'run_models.json').read_text());entry=next(e for e in entries if e['id']==args.model)
    manifest=json.loads((ROOT/'data/manifest.json').read_text())
    holdout=json.loads((ROOT/'audit/holdout.json').read_text());assert holdout['complete']
    adapter=Adapter(entry)
    dump(ROOT/'results'/entry['id'].replace('/','--')/'runtime.json',adapter.runtime)
    for name,meta in manifest['datasets'].items():
        if args.datasets and name not in args.datasets:continue
        docs=json.loads((ROOT/f'data/{name}.panel.json').read_text())
        for track in ['benchmark_schema','familiar_training_schema']:
            if args.tracks and track not in args.tracks:continue
            try:evaluate(adapter,name,meta,track,docs,holdout['matches'],args.deadline)
            except TimeoutError:raise
            except Exception as exc:
                dump(ROOT/'results'/entry['id'].replace('/','--')/f'{name}.{track}.error.json',{'error':repr(exc),'traceback':traceback.format_exc()})
                print('CELL_ERROR',name,track,repr(exc),flush=True)
    print('MODEL_COMPLETE',entry['id'],flush=True)

if __name__=='__main__':main()