Skip to content

Commit 26ae28d

Browse files
authored
Merge pull request #21 from Agentiix/fix-examples-logging
2 parents 357d344 + 5f1d684 commit 26ae28d

12 files changed

Lines changed: 217 additions & 38 deletions

File tree

agentix/log/__init__.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@
1919

2020
import logging
2121

22-
__all__ = ["install_worker_bridge"]
22+
from agentix.log._config import configure_logging
23+
24+
__all__ = ["configure_logging", "install_worker_bridge"]
2325

2426

2527
def install_worker_bridge(level: int = logging.NOTSET) -> logging.Handler:

agentix/log/_bridge.py

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,13 +2,13 @@
22

33
from __future__ import annotations
44

5-
import asyncio
65
import logging
76
from typing import Any
87

98
import socketio
109

1110
from agentix import sio as _sio
11+
from agentix.log._config import LOG_CONTEXT_ATTR
1212

1313
NAMESPACE = "/log"
1414

@@ -48,10 +48,7 @@ def emit(self, record: logging.LogRecord) -> None:
4848
try:
4949
payload = _record_payload(record)
5050
ns = _get_worker_namespace()
51-
asyncio.get_running_loop().create_task(ns.emit("record", payload))
52-
except RuntimeError:
53-
# No running loop — drop the record (worker is between async ticks).
54-
pass
51+
_sio._emit_nowait(ns.namespace, "record", payload)
5552
except Exception:
5653
self.handleError(record)
5754

@@ -86,6 +83,7 @@ def emit(self, record: logging.LogRecord) -> None:
8683
# unconditionally so a record produced on 3.12+ doesn't try to
8784
# smuggle `taskName` through `extra=` into a fresh record.
8885
"taskName",
86+
LOG_CONTEXT_ATTR,
8987
}
9088
)
9189

@@ -105,6 +103,7 @@ def _record_payload(record: logging.LogRecord) -> dict[str, Any]:
105103
"exc_text": record.exc_text
106104
or (logging.Formatter().formatException(record.exc_info) if record.exc_info else None),
107105
"stack_info": record.stack_info,
106+
LOG_CONTEXT_ATTR: getattr(record, LOG_CONTEXT_ATTR, None),
108107
"extras": extras or None,
109108
}
110109

@@ -166,6 +165,8 @@ def _replay_record(payload: dict[str, Any]) -> None:
166165
record.exc_text = str(payload["exc_text"])
167166
if payload.get("stack_info"):
168167
record.stack_info = str(payload["stack_info"])
168+
if payload.get(LOG_CONTEXT_ATTR):
169+
setattr(record, LOG_CONTEXT_ATTR, str(payload[LOG_CONTEXT_ATTR]))
169170
logger.handle(record)
170171

171172

agentix/log/_config.py

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,100 @@
1+
"""Small env-driven logging setup shared by host, runtime, and worker."""
2+
3+
from __future__ import annotations
4+
5+
import logging
6+
import os
7+
import socket
8+
import sys
9+
from typing import TextIO
10+
11+
LOG_CONTEXT_ATTR = "agentix_context"
12+
DEFAULT_LOG_FORMAT = f"%(asctime)s [%({LOG_CONTEXT_ATTR})s] [%(name)s] %(message)s"
13+
14+
_context = "host"
15+
_factory_installed = False
16+
_previous_factory = logging.getLogRecordFactory()
17+
18+
19+
class _SafeFormatMap(dict[str, str]):
20+
def __missing__(self, key: str) -> str:
21+
return "{" + key + "}"
22+
23+
24+
def configure_logging(
25+
*,
26+
default_context: str = "host",
27+
stream: TextIO | None = None,
28+
force: bool = False,
29+
) -> None:
30+
"""Configure stdlib logging from Agentix env vars.
31+
32+
Env vars:
33+
AGENTIX_LOG_LEVEL: logging level name, default INFO.
34+
AGENTIX_LOG_FORMAT: stdlib %-style format string.
35+
AGENTIX_LOG_CONTEXT: context template for this process.
36+
37+
Context templates support `{uname}`, `{hostname}`, `{pid}`, and
38+
`{id}`. The worker spawner sets `AGENTIX_WORKER_ID`, which backs
39+
`{id}` inside worker processes.
40+
"""
41+
set_log_context(os.environ.get("AGENTIX_LOG_CONTEXT", default_context))
42+
level_name = os.environ.get("AGENTIX_LOG_LEVEL", "INFO").upper()
43+
level = getattr(logging, level_name, logging.INFO)
44+
log_format = os.environ.get("AGENTIX_LOG_FORMAT", DEFAULT_LOG_FORMAT)
45+
logging.basicConfig(
46+
level=level,
47+
stream=stream or sys.stderr,
48+
format=log_format,
49+
force=force,
50+
)
51+
52+
53+
def set_log_context(template: str) -> None:
54+
global _context
55+
_context = _expand_context(template)
56+
_install_record_factory()
57+
58+
59+
def get_log_context() -> str:
60+
return _context
61+
62+
63+
def _expand_context(template: str) -> str:
64+
hostname = socket.gethostname()
65+
values = _SafeFormatMap(
66+
{
67+
"hostname": hostname,
68+
"uname": hostname,
69+
"pid": str(os.getpid()),
70+
"id": os.environ.get("AGENTIX_WORKER_ID", str(os.getpid())),
71+
}
72+
)
73+
try:
74+
return template.format_map(values)
75+
except ValueError:
76+
return template
77+
78+
79+
def _install_record_factory() -> None:
80+
global _factory_installed
81+
if _factory_installed:
82+
return
83+
84+
def record_factory(*args, **kwargs):
85+
record = _previous_factory(*args, **kwargs)
86+
if not hasattr(record, LOG_CONTEXT_ATTR):
87+
setattr(record, LOG_CONTEXT_ATTR, _context)
88+
return record
89+
90+
logging.setLogRecordFactory(record_factory)
91+
_factory_installed = True
92+
93+
94+
__all__ = [
95+
"DEFAULT_LOG_FORMAT",
96+
"LOG_CONTEXT_ATTR",
97+
"configure_logging",
98+
"get_log_context",
99+
"set_log_context",
100+
]

agentix/runtime/server/app.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,12 +23,13 @@
2323
from fastapi import FastAPI
2424

2525
from agentix import __version__
26+
from agentix.log import configure_logging
2627
from agentix.runtime.server.sio import make_sio
2728
from agentix.runtime.server.worker import RuntimeWorkerClient
2829
from agentix.runtime.shared.models import HealthResponse
2930

31+
configure_logging(default_context="sandbox-{uname}")
3032
logger = logging.getLogger("agentix.runtime")
31-
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(name)s] %(message)s")
3233

3334

3435
@asynccontextmanager

agentix/runtime/server/sio.py

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -190,9 +190,7 @@ async def on_cancel(sid: str, data: Any) -> None:
190190
# `sio_emit` — worker wants to broadcast an event on a namespace;
191191
# we pack the payload and call sio.emit there.
192192

193-
_broadcast_tasks: set[asyncio.Task] = set()
194-
195-
def _on_worker_sio_frame(frame: dict[str, Any]) -> None:
193+
async def _on_worker_sio_frame(frame: dict[str, Any]) -> None:
196194
kind = frame.get("type")
197195
namespace = frame.get("namespace")
198196
if not isinstance(namespace, str) or not namespace.startswith("/"):
@@ -201,11 +199,7 @@ def _on_worker_sio_frame(frame: dict[str, Any]) -> None:
201199
event = frame.get("event")
202200
if not isinstance(event, str):
203201
return
204-
task = asyncio.create_task(
205-
sio.emit(event, pack(frame.get("data")), namespace=namespace),
206-
)
207-
_broadcast_tasks.add(task)
208-
task.add_done_callback(_broadcast_tasks.discard)
202+
await sio.emit(event, pack(frame.get("data")), namespace=namespace)
209203
elif kind == "sio_open":
210204
if namespace in opened_namespaces or namespace == "/":
211205
return

agentix/runtime/server/worker/client.py

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010

1111
import asyncio
1212
import contextlib
13+
import inspect
1314
import logging
1415
import os
1516
import sys
@@ -128,7 +129,7 @@ async def send_inbound(self, namespace: str, event: str, data: Any) -> None: ...
128129
async def shutdown(self) -> None: ...
129130

130131

131-
SioFrameHandler = Callable[[dict[str, Any]], None]
132+
SioFrameHandler = Callable[[dict[str, Any]], Any]
132133
"""Called from the worker read loop for each `sio_emit` or
133134
`sio_subscribe` frame the worker produces. The SIO server layer
134135
installs one; the in-process backend has no transport hop."""
@@ -173,6 +174,7 @@ def __init__(
173174
) -> None:
174175
self._python = python
175176
self._runtime_bin_dir = runtime_bin_dir
177+
self._worker_id = _new_id()[:8]
176178

177179
self._proc: asyncio.subprocess.Process | None = None
178180
self._send_lock = asyncio.Lock()
@@ -188,6 +190,13 @@ def __init__(
188190
async def start(self) -> None:
189191
env = _clean_worker_env(self._runtime_bin_dir)
190192
env["AGENTIX_WORKER_IMPORT_ROOT"] = str(_WORKER_IMPORT_ROOT)
193+
env["AGENTIX_WORKER_ID"] = self._worker_id
194+
env["AGENTIX_LOG_CONTEXT"] = env.get(
195+
"AGENTIX_WORKER_LOG_CONTEXT",
196+
"sandbox-{uname}-worker-{id}",
197+
)
198+
if worker_log_format := env.get("AGENTIX_WORKER_LOG_FORMAT"):
199+
env["AGENTIX_LOG_FORMAT"] = worker_log_format
191200
self._proc = await asyncio.create_subprocess_exec(
192201
self._python,
193202
"-c",
@@ -235,7 +244,7 @@ async def _read_loop(self) -> None:
235244
frame = await read_frame(self._proc.stdout)
236245
if frame is None:
237246
break
238-
self._on_frame(frame)
247+
await self._on_frame(frame)
239248
except Exception:
240249
logger.exception("runtime worker read loop crashed")
241250
finally:
@@ -246,7 +255,7 @@ async def _read_loop(self) -> None:
246255
fut.set_exception(RuntimeError(err.message))
247256
self._pending.clear()
248257

249-
def _on_frame(self, frame: dict[str, Any]) -> None:
258+
async def _on_frame(self, frame: dict[str, Any]) -> None:
250259
kind = frame.get("type")
251260
if kind == "ready":
252261
self._ready.set()
@@ -268,7 +277,9 @@ def _on_frame(self, frame: dict[str, Any]) -> None:
268277
elif kind in ("sio_emit", "sio_open"):
269278
if self._sio_handler is not None:
270279
try:
271-
self._sio_handler(frame)
280+
result = self._sio_handler(frame)
281+
if inspect.isawaitable(result):
282+
await result
272283
except Exception:
273284
logger.debug("sio frame handler raised; dropping", exc_info=True)
274285
else:

agentix/runtime/server/worker/process.py

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -211,12 +211,9 @@ async def _amain() -> None:
211211

212212

213213
def main() -> None:
214-
level_name = os.environ.get("AGENTIX_LOG_LEVEL", "INFO").upper()
215-
level = getattr(logging, level_name, logging.INFO)
216-
logging.basicConfig(
217-
level=level,
214+
_log.configure_logging(
215+
default_context="sandbox-{uname}-worker-{id}",
218216
stream=sys.stderr,
219-
format="%(asctime)s [%(name)s] %(message)s",
220217
)
221218
try:
222219
asyncio.run(_amain())

agentix/sio.py

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -154,14 +154,7 @@ def _auto_register(self) -> None:
154154

155155
async def emit(self, event: str, data: Any = None) -> None:
156156
"""Emit `event` on this namespace to all connected hosts."""
157-
_bridge.send_frame(
158-
{
159-
"type": "sio_emit",
160-
"namespace": self.namespace,
161-
"event": event,
162-
"data": data,
163-
}
164-
)
157+
_emit_nowait(self.namespace, event, data)
165158

166159
def on(self, event: str, handler: Handler) -> None:
167160
"""Register an additional handler for `event`."""
@@ -291,6 +284,17 @@ def _is_installed() -> bool:
291284
return _bridge.is_installed()
292285

293286

287+
def _emit_nowait(namespace: str, event: str, data: Any = None) -> None:
288+
_bridge.send_frame(
289+
{
290+
"type": "sio_emit",
291+
"namespace": namespace,
292+
"event": event,
293+
"data": data,
294+
}
295+
)
296+
297+
294298
__all__ = [
295299
"RESERVED_NAMESPACES",
296300
"Namespace",

agentix/trace/_bridge.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,6 @@
1111

1212
from __future__ import annotations
1313

14-
import asyncio
1514
import contextvars
1615
from typing import Any
1716

@@ -76,7 +75,7 @@ def _emit(self, event: str, payload: dict[str, Any]) -> None:
7675
if not _sio._is_installed():
7776
return
7877
try:
79-
asyncio.get_running_loop().create_task(self._ns.emit(event, payload))
78+
_sio._emit_nowait(self._ns.namespace, event, payload)
8079
except RuntimeError:
8180
pass
8281

examples/eval-cc-swe/runner.py

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@
4040

4141
from agentix import RuntimeClient
4242
from agentix.deployment.base import SandboxConfig, session
43+
from agentix.log import configure_logging
4344

4445
WORKDIR = "/testbed"
4546
logger = logging.getLogger("eval_cc_swe.runner")
@@ -419,10 +420,7 @@ async def main(argv: list[str] | None = None) -> int:
419420
print("error: --concurrency must be >= 1", file=sys.stderr)
420421
return 2
421422

422-
logging.basicConfig(
423-
level=logging.INFO,
424-
format="%(asctime)s [%(name)s] %(message)s",
425-
)
423+
configure_logging(default_context="host")
426424

427425
ds = _load_instances_dataset(args.dataset, split=args.split, dataset_file=args.dataset_file)
428426
instances = _selected_instances(

0 commit comments

Comments
 (0)