Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@

A `PtyRegistry` adds named get-or-create sessions, injected shell setup (`rc`), and inactivity culling. It carries no server, no framework, and no session persistence beyond the process. Web exposure belongs to the embedding app (e.g. [jupygate](https://github.com/AnswerDotAI/jupygate)), and durable session management belongs to tmux.

`ptymini.bg` is the sync surface, for callers outside an event loop: kernel tools, plain scripts, an LLM deciding between actions. It is the [bgterm](https://github.com/AnswerDotAI/bgterm) API, folded in, with sid-named sessions and blocking calls. The wait parameters are `fastmux.bg`’s. `wait_ms` bounds the wait for new output. `until=` returns as soon as the accumulated text matches that regex. `settle_ms` keeps collecting until output has stopped for that long.

### Installation

``` sh
Expand Down
32 changes: 27 additions & 5 deletions nbs/01_bg.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
"id": "3a336167",
"metadata": {},
"source": [
"The [bgterm](https://github.com/AnswerDotAI/bgterm) package, folded in: a *sync*, sid-based API for background terminal sessions — start once, send input later, wait a bounded time, read what arrived since your last look. It is the cursor view of a `PtySession`, for callers that live outside an event loop (kernel tools, plain scripts, an LLM deciding between actions): every session shares one private asyncio loop on a daemon thread, and each call marshals onto it and blocks, so `write_stdin(sid, \"2+2\\n\", 500)` means what it always meant. The buffering, paging, and drop accounting are the `Ring`'s, which was extracted from bgterm in the first place; `PollResult` and every function signature are unchanged from the original package."
"The [bgterm](https://github.com/AnswerDotAI/bgterm) package, folded in: a *sync*, sid-based API for background terminal sessions — start once, send input later, wait a bounded time, read what arrived since your last look. It is the cursor view of a `PtySession`, for callers that live outside an event loop (kernel tools, plain scripts, an LLM deciding between actions): every session shares one private asyncio loop on a daemon thread, and each call marshals onto it and blocks, so `write_stdin(sid, \"2+2\\n\", 500)` means what it always meant. The buffering, paging, and drop accounting are the `Ring`'s, which was extracted from bgterm in the first place. The wait parameters are `fastmux.bg`'s. `wait_ms` bounds the wait for new output. `until=` returns as soon as the accumulated text matches that regex. `settle_ms` keeps collecting until output has stopped for that long."
]
},
{
Expand Down Expand Up @@ -75,7 +75,7 @@
"id": "2e0e9aa4",
"metadata": {},
"source": [
"The sid table and functional API, unchanged from bgterm — the sid is how a kernel tool names a session across calls:"
"The sid table and functional API — the sid is how a kernel tool names a session across calls. The wait parameters follow `fastmux.bg`, with two differences the substrate forces. Waits are event-driven: a condvar wakes on new bytes or death. There is no `interval_ms`. Reading consumes the stream, and `until` can only match output accumulated within the current call. Match on text only the awaited output can produce, such as a fresh prompt. The whole call is bounded by `wait_ms + settle_ms`."
]
},
{
Expand Down Expand Up @@ -191,7 +191,7 @@
"id": "466e213b",
"metadata": {},
"source": [
"The whole flow, as a reader would use it — start a REPL-ish child, poll for its banner, send input, read what arrived. Everything here is a plain sync call:"
"The whole flow, as a reader would use it — start a REPL-ish child, wait for its banner, then send input and name the reply to wait for. `until` replaces the write-then-poll-again loop with one call:"
]
},
{
Expand All @@ -205,13 +205,35 @@
" \"print('ready', flush=True)\\nimport sys\\nfor line in sys.stdin: print(f'ACK:{line.strip()}', flush=True)\"])\n",
"first = poll(sid, 5000)\n",
"assert 'ready' in first.text\n",
"r = write_stdin(sid, 'hello\\n', 500)\n",
"if 'ACK:hello' not in r.text: r = poll(sid, 2000)\n",
"r = write_stdin(sid, 'hello\\n', 2000, until=r'ACK:hello')\n",
"assert 'ACK:hello' in r.text\n",
"assert sid in list_sessions()\n",
"r.text"
]
},
{
"cell_type": "markdown",
"id": "d7c6ab31",
"metadata": {},
"source": [
"`settle_ms` is for replies whose shape you don't know. After the wait, the call keeps collecting until the stream has been quiet for that long. A multi-burst reply comes back in one call. A child that dies ends the settle at once."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "c93a1558",
"metadata": {},
"outputs": [],
"source": [
"sid2 = start_bgterm([sys.executable, '-u', '-c',\n",
" \"import time\\nprint('part one', flush=True)\\ntime.sleep(0.2)\\nprint('part two', flush=True)\"])\n",
"r = poll(sid2, 5000, settle_ms=500)\n",
"assert 'part one' in r.text and 'part two' in r.text\n",
"close_bgterm(sid2)\n",
"r.text"
]
},
{
"cell_type": "markdown",
"id": "85d1958a",
Expand Down
8 changes: 8 additions & 0 deletions nbs/index.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@
"A `PtyRegistry` adds named get-or-create sessions, injected shell setup (`rc`), and inactivity culling. It carries no server, no framework, and no session persistence beyond the process. Web exposure belongs to the embedding app (e.g. [jupygate](https://github.com/AnswerDotAI/jupygate)), and durable session management belongs to tmux."
]
},
{
"cell_type": "markdown",
"id": "5802ae72",
"metadata": {},
"source": [
"`ptymini.bg` is the sync surface, for callers outside an event loop: kernel tools, plain scripts, an LLM deciding between actions. It is the [bgterm](https://github.com/AnswerDotAI/bgterm) API, folded in, with sid-named sessions and blocking calls. The wait parameters are `fastmux.bg`'s. `wait_ms` bounds the wait for new output. `until=` returns as soon as the accumulated text matches that regex. `settle_ms` keeps collecting until output has stopped for that long."
]
},
{
"cell_type": "markdown",
"id": "f5cf0a39",
Expand Down
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,9 @@ dependencies = ['fastcore>=2.1.18']
[project.optional-dependencies]
dev = ["maturin>=1.0,<2.0", "pytest", "nbdev>=3", "fastship"]

[project.entry-points.pyskills]
ptymini = "ptymini.skill"

[project.urls]
Repository = "https://github.com/AnswerDotAI/ptymini"
Documentation = "https://AnswerDotAI.github.io/ptymini/"
Expand Down
62 changes: 43 additions & 19 deletions python/ptymini/bg.py
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
"""The bgterm API: sync, cursor-paged background terminal sessions

The [bgterm](https://github.com/AnswerDotAI/bgterm) package, folded in: a *sync*, sid-based API for background terminal sessions — start once, send input later, wait a bounded time, read what arrived since your last look. It is the cursor view of a pty session, for callers that live outside an event loop (kernel tools, plain scripts, an LLM deciding between actions): each call runs directly against the sync Rust core and blocks, with the GIL released while waiting, so `write_stdin(sid, "2+2\n", 500)` means what it always meant. The buffering, paging, and drop accounting are the `Ring`'s, which was extracted from bgterm in the first place; `PollResult` and every function signature are unchanged from the original package.
The [bgterm](https://github.com/AnswerDotAI/bgterm) package, folded in: a *sync*, sid-based API for background terminal sessions — start once, send input later, wait a bounded time, read what arrived since your last look. It is the cursor view of a pty session, for callers that live outside an event loop (kernel tools, plain scripts, an LLM deciding between actions): each call runs directly against the sync Rust core and blocks, with the GIL released while waiting, so `write_stdin(sid, "2+2\n", 500)` means what it always meant. The buffering, paging, and drop accounting are the `Ring`'s, which was extracted from bgterm in the first place. The wait parameters are `fastmux.bg`'s. `wait_ms` bounds the wait for new output. `until=` returns as soon as the accumulated text matches that regex. `settle_ms` keeps collecting until output has stopped for that long.

Docs: https://AnswerDotAI.github.io/ptymini/bg.html.md"""


__all__ = ['DEFAULT_MAX_BUFFER_BYTES', 'DEFAULT_MAX_OUTPUT_BYTES', 'Cmd', 'BgtermError', 'PollResult', 'list_sessions',
'start_bgterm', 'write_stdin', 'poll', 'read', 'wait', 'terminate', 'kill', 'close_bgterm', 'Session']

import itertools, os, signal, threading
import itertools, os, re, signal, threading, time
from dataclasses import dataclass
from ._core import PtyCore

Expand Down Expand Up @@ -69,17 +69,33 @@ def _read(self, max_output_bytes):
return PollResult(payload.decode(self.encoding, errors=self.errors), payload, start_offset, end_offset,
r.start, r.end, len(payload), max(0, r.end - end_offset), dropped, self.s.alive, self.s.exit_code)

def _peek(self):
payload, _, _ = self.s.read_from(self.cursor, None)
return payload.decode(self.encoding, errors=self.errors)

def read(self, max_output_bytes=DEFAULT_MAX_OUTPUT_BYTES): return self._read(max_output_bytes)

def write_stdin(self, chars='', yield_time_ms=0, max_output_bytes=DEFAULT_MAX_OUTPUT_BYTES):
def write_stdin(self, chars='', wait_ms=0, max_output_bytes=DEFAULT_MAX_OUTPUT_BYTES, until=None, settle_ms=0):
if chars:
data = chars if isinstance(chars, bytes) else chars.encode(self.encoding)
try: self.s.write(data)
except OSError as e: raise BgtermError('session PTY is no longer writable') from e
if yield_time_ms > 0 and self.cursor >= self.s.end and self.s.alive: self.s.wait_change(timeout=yield_time_ms / 1000)
s, deadline = self.s, time.monotonic() + wait_ms / 1000
if until is not None:
while not re.search(until, self._peek()) and time.monotonic() < deadline:
seen = s.end
s.wait_change(seen, max(0.0, deadline - time.monotonic()))
if s.end == seen: break
elif wait_ms > 0 and self.cursor >= s.end and s.alive: s.wait_change(self.cursor, wait_ms / 1000)
if settle_ms:
hard = deadline + settle_ms / 1000
while time.monotonic() < hard:
seen = s.end
if not s.wait_change(seen, min(settle_ms / 1000, hard - time.monotonic())) or s.end == seen: break
return self._read(max_output_bytes)

def poll(self, yield_time_ms=0, max_output_bytes=DEFAULT_MAX_OUTPUT_BYTES): return self.write_stdin('', yield_time_ms, max_output_bytes)
def poll(self, wait_ms=0, max_output_bytes=DEFAULT_MAX_OUTPUT_BYTES, until=None, settle_ms=0):
return self.write_stdin('', wait_ms, max_output_bytes, until, settle_ms)

def wait(self, timeout_ms=None): return self.s.wait(timeout=None if timeout_ms is None else max(timeout_ms, 0) / 1000)

Expand Down Expand Up @@ -135,22 +151,26 @@ def start_bgterm(


def write_stdin(
sid:int, # Session id from `start_bgterm`
chars:str|bytes='', # Input to write; str encodes with the session's encoding, bytes pass through
yield_time_ms:int=0, # Max ms to wait for output when none is unread
sid:int, # Session id from `start_bgterm`
chars:str|bytes='', # Input to write; str encodes with the session's encoding, bytes pass through
wait_ms:int=0, # Max ms to wait for output when none is unread (with `until`: for the match)
max_output_bytes:int=DEFAULT_MAX_OUTPUT_BYTES, # Page size cap on returned bytes; None returns everything unread
until:str=None, # Regex; return as soon as this call's accumulated unread text matches
settle_ms:int=0, # Then keep collecting until output stops for this long (at most `settle_ms` past the deadline)
):
"Write to a session PTY, wait briefly, and return unread output."
return _lookup(sid).write_stdin(chars, yield_time_ms, max_output_bytes)
"Write to a session PTY, wait as `poll` does, and return unread output."
return _lookup(sid).write_stdin(chars, wait_ms, max_output_bytes, until, settle_ms)


def poll(
sid:int, # Session id from `start_bgterm`
yield_time_ms:int=0, # Max ms to wait for output when none is unread
sid:int, # Session id from `start_bgterm`
wait_ms:int=0, # Max ms to wait for output when none is unread (with `until`: for the match)
max_output_bytes:int=DEFAULT_MAX_OUTPUT_BYTES, # Page size cap on returned bytes; None returns everything unread
until:str=None, # Regex; return as soon as this call's accumulated unread text matches
settle_ms:int=0, # Then keep collecting until output stops for this long (at most `settle_ms` past the deadline)
):
"Wait for unread session output, then return it without writing input."
return _lookup(sid).poll(yield_time_ms, max_output_bytes)
"Wait for unread session output (or an `until` match, or settling), then return it without writing input."
return _lookup(sid).poll(wait_ms, max_output_bytes, until, settle_ms)


def read(
Expand Down Expand Up @@ -243,18 +263,22 @@ def exit_code(self):

def write_stdin(self,
chars:str|bytes='', # Input to write; str encodes with the session's encoding, bytes pass through
yield_time_ms:int=0, # Max ms to wait for output when none is unread
wait_ms:int=0, # Max ms to wait for output when none is unread (with `until`: for the match)
max_output_bytes:int=DEFAULT_MAX_OUTPUT_BYTES, # Page size cap on returned bytes; None returns everything unread
until:str=None, # Regex; return as soon as this call's accumulated unread text matches
settle_ms:int=0, # Then keep collecting until output stops for this long (at most `settle_ms` past the deadline)
):
"Write to this session PTY, wait briefly, and return unread output."
return write_stdin(self.sid, chars, yield_time_ms, max_output_bytes)
"Write to this session PTY, wait as `poll` does, and return unread output."
return write_stdin(self.sid, chars, wait_ms, max_output_bytes, until, settle_ms)

def poll(self,
yield_time_ms:int=0, # Max ms to wait for output when none is unread
wait_ms:int=0, # Max ms to wait for output when none is unread (with `until`: for the match)
max_output_bytes:int=DEFAULT_MAX_OUTPUT_BYTES, # Page size cap on returned bytes; None returns everything unread
until:str=None, # Regex; return as soon as this call's accumulated unread text matches
settle_ms:int=0, # Then keep collecting until output stops for this long (at most `settle_ms` past the deadline)
):
"Wait for unread output from this session, then return it."
return poll(self.sid, yield_time_ms, max_output_bytes)
return poll(self.sid, wait_ms, max_output_bytes, until, settle_ms)

def read(self,
max_output_bytes:int=DEFAULT_MAX_OUTPUT_BYTES, # Page size cap on returned bytes; None returns everything unread
Expand Down
20 changes: 20 additions & 0 deletions python/ptymini/skill.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
r"""Drive line-oriented interactive programs in the background: REPLs, shells, ssh, debuggers. Start a pty session once, send input later, wait a bounded time for the reply, and read what arrived since your last look. Use this for CLI work. Use `fastmux` for TUIs, rich terminal apps, and terminals shared with the user.

The API is `ptymini.bg`, the bgterm interface. Calls are sync and sid-based. They run directly against the Rust pty core and need no event loop. Sessions live inside this process and end with it. Nothing is shared server-side.

The whole flow:

sid = start_bgterm(['ipython', '--simple-prompt'])
poll(sid, 5000, until=r'In \[1\]')
r = write_stdin(sid, '2+2\n', 5000, until=r'In \[2\]')
close_bgterm(sid)

The wait parameters are `fastmux.bg`'s. `wait_ms` bounds the wait for new output. `until=` returns as soon as the accumulated text matches that regex. Match on text only the awaited reply can produce, such as the next prompt. `settle_ms` keeps collecting until output has stopped for that long, and runs at most `settle_ms` past the `wait_ms` deadline. Every wait is bounded. On timeout the call returns whatever arrived. Two behaviors differ from fastmux. Waits are event-driven, and there is no `interval_ms`. Reading consumes the stream, and `until` never sees output an earlier call returned.

Every read-shaped call returns a `PollResult`. `text` and `data` hold the output. `remaining_bytes` and `truncated` report paging. `dropped_bytes` counts unread output lost when the ring overran its bound. `running` and `exit_code` report liveness; a negative exit code is the terminating signal number. A dead child ends any wait at once. `read()` returns immediately. `wait()` blocks until exit. `terminate`, `kill`, and `close_bgterm` end the child. `list_sessions()` lists live sids. `Session` wraps a sid for `with` blocks.

`ptymini.core` is the asyncio surface: multi-client attach, replay buffers, and a named registry. Run `doc(func)` for full parameter docments before first use.
"""

from .bg import *
from .bg import __all__
3 changes: 1 addition & 2 deletions tests/test_bg.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,7 @@ def test_functional_session_accepts_input_and_returns_echoed_output():
ready = poll(sid, 500)
assert "ready" in ready.text

reply = write_stdin(sid, "hello\n", 100)
if "ACK:hello" not in reply.text: reply = poll(sid, 500)
reply = write_stdin(sid, "hello\n", 2000, until=r"ACK:hello")
assert "ACK:hello" in reply.text
finally: close_bgterm(sid)

Expand Down