File size: 8,810 Bytes
9abace2
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
"""

Stdio Bridge Server

Part of SOVEREIGN PYTHON LLM ENGINE



Stdio-based bridge for C frontend communication.

Uses JSON-RPC 2.0 over stdin/stdout.

"""

import sys
import asyncio
import json
from typing import Any
from datetime import datetime

from ..agents.react import ReActAgent, ReActConfig
from ..tools.registry import ToolRegistry
from ..tools.approval import ApprovalEngine
from ..core.evidence import WORMLedger
from ..core.crypto import generate_signing_key
from ..runtime.providers.multi import MultiProvider
from ..models.entities import Task, Message, MessageRole
from pathlib import Path


class StdioBridge:
    """

    Stdio bridge server.



    Protocol:

    - C frontend writes JSON-RPC requests to our stdin

    - We write JSON-RPC responses to stdout

    - All messages are newline-delimited JSON

    """

    def __init__(self):
        # Initialize components
        signing_key = generate_signing_key()
        self.ledger = WORMLedger(Path("bridge_evidence.worm"), signing_key)
        self.registry = ToolRegistry()
        self.approval = ApprovalEngine(self.ledger)
        self.model = MultiProvider()  # OpenRouter → Ollama fallback

        # Initialize agent
        config = ReActConfig(max_steps=15, log_to_worm=True)
        self.agent = ReActAgent(
            self.model,
            self.registry,
            self.approval,
            self.ledger,
            config
        )

        # Request handlers
        self.handlers = {
            "agent.run": self._handle_agent_run,
            "agent.stream": self._handle_agent_stream,
            "tool.execute": self._handle_tool_execute,
            "tool.list": self._handle_tool_list,
            "chat.message": self._handle_chat_message,
            "health": self._handle_health,
        }

    async def run(self):
        """Run stdio bridge server"""
        # Log startup
        await self.ledger.append({
            "event": "bridge_startup",
            "timestamp": datetime.utcnow().isoformat()
        })

        # Write ready signal
        self._write_response({
            "jsonrpc": "2.0",
            "method": "bridge.ready",
            "params": {
                "version": "1.0.0",
                "capabilities": ["agent.run", "agent.stream", "tool.execute", "chat.message"]
            }
        })

        # Read loop
        while True:
            try:
                # Read line from stdin
                line = await asyncio.get_event_loop().run_in_executor(
                    None,
                    sys.stdin.readline
                )

                if not line:
                    break  # EOF

                line = line.strip()
                if not line:
                    continue

                # Parse request
                request = json.loads(line)

                # Handle request
                response = await self._handle_request(request)

                # Write response
                if response:
                    self._write_response(response)

            except KeyboardInterrupt:
                break
            except Exception as e:
                # Write error response
                self._write_response({
                    "jsonrpc": "2.0",
                    "error": {
                        "code": -32603,
                        "message": "Internal error",
                        "data": str(e)
                    },
                    "id": None
                })

        # Log shutdown
        await self.ledger.append({
            "event": "bridge_shutdown",
            "timestamp": datetime.utcnow().isoformat()
        })

    async def _handle_request(self, request: dict) -> dict | None:
        """Handle JSON-RPC request"""
        # Validate
        if request.get("jsonrpc") != "2.0":
            return {
                "jsonrpc": "2.0",
                "error": {
                    "code": -32600,
                    "message": "Invalid Request"
                },
                "id": request.get("id")
            }

        method = request.get("method")
        if not method:
            return {
                "jsonrpc": "2.0",
                "error": {
                    "code": -32600,
                    "message": "Method required"
                },
                "id": request.get("id")
            }

        params = request.get("params", {})
        request_id = request.get("id")

        # Find handler
        handler = self.handlers.get(method)
        if not handler:
            return {
                "jsonrpc": "2.0",
                "error": {
                    "code": -32601,
                    "message": f"Method not found: {method}"
                },
                "id": request_id
            }

        # Execute
        try:
            result = await handler(params)

            # Don't send response for notifications
            if request_id is None:
                return None

            return {
                "jsonrpc": "2.0",
                "result": result,
                "id": request_id
            }

        except Exception as e:
            return {
                "jsonrpc": "2.0",
                "error": {
                    "code": -32000,
                    "message": "Handler error",
                    "data": str(e)
                },
                "id": request_id
            }

    async def _handle_agent_run(self, params: dict) -> dict:
        """Handle agent.run request"""
        task_description = params.get("task")
        if not task_description:
            raise ValueError("task parameter required")

        # Create task
        task = Task(
            task_id=f"bridge-{datetime.utcnow().timestamp()}",
            description=task_description
        )

        # Run agent
        result = await self.agent.run(task)

        return {
            "result": result,
            "task_id": task.task_id
        }

    async def _handle_agent_stream(self, params: dict) -> dict:
        """Handle agent.stream request"""
        # Streaming would require SSE or WebSocket
        # For now, fall back to non-streaming
        return await self._handle_agent_run(params)

    async def _handle_tool_execute(self, params: dict) -> dict:
        """Handle tool.execute request"""
        tool_name = params.get("tool")
        tool_args = params.get("args", {})

        if not tool_name:
            raise ValueError("tool parameter required")

        # Get tool
        tool = self.registry.get(tool_name)
        if not tool:
            raise ValueError(f"Tool not found: {tool_name}")

        # Execute
        result = await tool.handler(**tool_args)

        return {
            "tool": tool_name,
            "result": str(result)
        }

    async def _handle_tool_list(self, params: dict) -> dict:
        """Handle tool.list request"""
        tools = self.registry.list_all()

        return {
            "tools": [
                {
                    "id": tool.tool_id,
                    "description": tool.description,
                    "risk_class": tool.risk_class.value
                }
                for tool in tools[:50]  # Limit to 50
            ]
        }

    async def _handle_chat_message(self, params: dict) -> dict:
        """Handle chat.message request"""
        message_text = params.get("message")
        if not message_text:
            raise ValueError("message parameter required")

        # Generate response using model
        messages = [
            {
                "role": "user",
                "content": message_text
            }
        ]

        response = await self.model.invoke_model(
            model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
            messages=messages,
            max_tokens=2048
        )

        # Extract text
        reply = response["content"][0]["text"]

        return {
            "reply": reply
        }

    async def _handle_health(self, params: dict) -> dict:
        """Handle health check"""
        return {
            "status": "ok",
            "timestamp": datetime.utcnow().isoformat()
        }

    def _write_response(self, response: dict) -> None:
        """Write JSON response to stdout"""
        json_str = json.dumps(response)
        sys.stdout.write(json_str + "\n")
        sys.stdout.flush()


def main():
    """Entry point for stdio bridge"""
    bridge = StdioBridge()
    asyncio.run(bridge.run())


if __name__ == "__main__":
    main()