File size: 6,293 Bytes
b4b0f75
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""JSONL and local HTTP entry points sharing the supported decision engine."""
import argparse
from contextlib import contextmanager
from datetime import datetime,timezone
from http.server import BaseHTTPRequestHandler,ThreadingHTTPServer
import json
from pathlib import Path
import signal,sys,threading
import urllib.error
from .backend import Client,NativeServer
from .decision import DecisionEngine

CAPABILITIES=dict(api_version='1',input_fields=['state','question','options'],min_options=2,max_options=26,
                  option_ids='opaque',numeric_values='parsed from option descriptions only',reasoning=False,
                  modalities=['text'],agent_execution='experimental; not exposed by this HTTP service',
                  compatibility='This is a decision API, not complete Jev/OpenJev API compatibility.')


def server(engine,host='0.0.0.0',port=9304):
    slots=threading.BoundedSemaphore(2)
    class Handler(BaseHTTPRequestHandler):
        def log_message(self,*args):pass
        def send_json(self,status,payload):
            data=json.dumps(payload,ensure_ascii=False).encode()
            self.send_response(status);self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(data)));self.end_headers()
            try:self.wfile.write(data)
            except (BrokenPipeError,ConnectionResetError):pass
        def do_GET(self):
            if self.path=='/health':
                try:
                    engine.client.health()
                    self.send_json(200,dict(status='ready',profile='base' if engine.client.adapter_id is None else 'adapter'))
                except (OSError,ValueError):self.send_json(503,dict(status='backend-unavailable'))
            elif self.path=='/v1/capabilities':self.send_json(200,CAPABILITIES)
            else:self.send_json(404,dict(error='Unknown endpoint'))
        def do_POST(self):
            if self.path!='/v1/decide':self.send_json(404,dict(error='Unknown endpoint'));return
            if not slots.acquire(blocking=False):self.send_json(503,dict(error='Decision worker and bounded queue are busy'));return
            try:
                self.connection.settimeout(125)
                length=int(self.headers.get('Content-Length','0'))
                if not 0<length<=65536:self.send_json(413,dict(error='Provide a JSON body of at most 64 KiB'));return
                data=self.rfile.read(length)
                if len(data)!=length:raise ValueError('Incomplete request body')
                request=json.loads(data)
                self.send_json(200,engine.decide(request))
            except (ValueError,UnicodeError,TypeError) as exc:self.send_json(400,dict(error=str(exc)))
            except (OSError,RuntimeError,KeyError,IndexError) as exc:self.send_json(502,dict(error='Backend decision failed',detail=str(exc)))
            finally:slots.release()
    return ThreadingHTTPServer((host,port),Handler)


def arguments(description):
    p=argparse.ArgumentParser(description=description)
    p.add_argument('--base',default='http://127.0.0.1:5991');p.add_argument('--profile',choices=['base','current'],default='base')
    p.add_argument('--adapter-id',type=int,default=0);p.add_argument('--launch',action='store_true')
    p.add_argument('--model',default='models/Ternary-Bonsai-2-27B-PQ2_0.gguf')
    p.add_argument('--adapter',help='Explicit adapter path required to launch profile current')
    p.add_argument('--binary',default='llama.cpp-b2/build/bin/llama-server');p.add_argument('--gpu',type=int,default=0)
    p.add_argument('--native-port',type=int,default=5991);p.add_argument('--run-dir')
    return p


@contextmanager
def engine_for(args):
    if args.launch:
        if args.profile=='current' and not args.adapter:raise ValueError('Current profile requires an explicit --adapter path; its weight license differs from the base')
        if args.profile=='base' and args.adapter:raise ValueError('Base profile does not load adapters')
        if args.adapter_id!=0:raise ValueError('Launching one adapter requires adapter ID 0')
        run=Path(args.run_dir) if args.run_dir else Path('runs')/datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S.%fZ')
        run.mkdir(parents=True,exist_ok=False)
        with NativeServer(run,args.gpu,args.native_port,args.model,args.adapter if args.profile=='current' else None,args.binary) as client:
            yield DecisionEngine(client)
    else:
        yield DecisionEngine(Client(args.base,None if args.profile=='base' else args.adapter_id))


@contextmanager
def graceful_termination():
    """Let service-manager termination unwind owned backend contexts."""
    def interrupt(signum,frame):raise KeyboardInterrupt
    previous=signal.signal(signal.SIGTERM,interrupt)
    try:yield
    finally:signal.signal(signal.SIGTERM,previous)


def cli():
    p=arguments('Bonsai text-choice decisions: JSONL input and JSONL output.');p.add_argument('--input');args=p.parse_args()
    try:
        with graceful_termination(),engine_for(args) as engine:
            stream=open(args.input) if args.input else sys.stdin
            try:
                for line in stream:
                    if not line.strip():continue
                    try:result=engine.decide(json.loads(line))
                    except (ValueError,TypeError,OSError,RuntimeError) as exc:result=dict(error=str(exc))
                    print(json.dumps(result),flush=True)
            finally:
                if args.input:stream.close()
    except KeyboardInterrupt:pass
    except (ValueError,RuntimeError) as exc:p.error(str(exc))


def http_cli():
    p=arguments('Serve the Bonsai decision API.');p.add_argument('--host',default='0.0.0.0');p.add_argument('--port',type=int,default=9304);args=p.parse_args()
    try:
        with graceful_termination(),engine_for(args) as engine:
            httpd=server(engine,args.host,args.port)
            print(json.dumps(dict(listening=f'http://{args.host}:{args.port}',profile=args.profile)),file=sys.stderr,flush=True)
            try:httpd.serve_forever(poll_interval=.2)
            except KeyboardInterrupt:pass
            finally:httpd.server_close()
    except KeyboardInterrupt:pass
    except (ValueError,RuntimeError) as exc:p.error(str(exc))


if __name__=='__main__':cli()