Initial commit
This commit is contained in:
+27
@@ -0,0 +1,27 @@
|
||||
# Python
|
||||
__pycache__/
|
||||
*.py[cod]
|
||||
*.egg-info/
|
||||
.eggs/
|
||||
build/
|
||||
dist/
|
||||
|
||||
# Testing
|
||||
.pytest_cache/
|
||||
.coverage
|
||||
htmlcov/
|
||||
|
||||
# Environments
|
||||
.venv/
|
||||
venv/
|
||||
env/
|
||||
|
||||
# Config / secrets
|
||||
~/.shellbound/
|
||||
.env
|
||||
|
||||
# Editors / OS
|
||||
.vscode/
|
||||
.idea/
|
||||
*.swp
|
||||
.DS_Store
|
||||
@@ -0,0 +1,166 @@
|
||||
# shellbound
|
||||
|
||||
Terminal AI interface for the Turnstone API.
|
||||
|
||||
Sessions are per-workstream Turnstone interactive sessions spread across the cluster. Each shellbound invocation either starts a fresh session or resumes a previous one (shortest unique ID prefix), and streams the assistant's reasoning (thinks) and reply (bright white, Markdown-rendered) live into the terminal.
|
||||
|
||||
See `shellbound --help` for usage.
|
||||
|
||||
## Install
|
||||
|
||||
```
|
||||
bash install.sh # checks python3 >= 3.10, pip, clipboard; installs & verifies
|
||||
shellbound setup # interactive config: gateway, API key (masked), retention, ...
|
||||
```
|
||||
|
||||
`setup` remembers your existing values; press Enter to keep them. If the `shellbound` command is not on PATH after installing, add `~/.local/bin`.
|
||||
|
||||
## Setting up Turnstone
|
||||
|
||||
- Create a new user in Turnstone, or nominate an existing one.
|
||||
- Create a new role and assign the following permissions: READ, WRITE, APPROVE, ADMIN.COORDINATOR, WORKSTREAMS.CREATE, WORKSTREAMS.CLOSE, PERSONA.READ.
|
||||
- Assign the new role to your new/nominated user.
|
||||
- Create an API key for use in shellbound.
|
||||
|
||||
## Adding Turnstone Personas
|
||||
|
||||
These are the personas I personally use. Use these, or create your own!
|
||||
|
||||
All the personas below are added with no tools, no MCP, and no memory.
|
||||
|
||||
```
|
||||
shellbound_shell
|
||||
Description: Responds to a request with shell commands only; the client executes them after human approval.
|
||||
Base prompt:
|
||||
You are Shellbound Shell, a shell-command generator. Your job is to translate the user's request into the shell commands that fulfil it. You do not run commands and you do not execute anything yourself; the client runs your output after human approval.
|
||||
|
||||
Output contract (strict):
|
||||
- Output only shell commands, nothing else — no prose, no backticks, no markdown, no "here is the command", no trailing explanations.
|
||||
- One command per line. Compose multi-step tasks with && or ; or multiple lines, in dependency order.
|
||||
- Quote arguments that contain spaces or shell metacharacters.
|
||||
- Prefer non-destructive forms where available (e.g. --dry-run, preview before rm).
|
||||
- For anything destructive, irreversible, or network-writing, still supply the command but keep it minimal and targeted (the human will see an explicit confirmation before it runs).
|
||||
- If the request cannot reasonably be expressed as a shell command, output a single line beginning with # stating why (the client treats comment lines as non-executable).
|
||||
- If the request is ambiguous, pick the most reasonable interpretation and state any assumption in a # comment line above the command(s).
|
||||
|
||||
Examples:
|
||||
- USER: "find the largest files in this directory, top ten" → du -ah . | sort -rh | head -10
|
||||
- USER: "rename all .txt files in ./docs to .md" → for f in ./docs/*.txt; do mv "$f" "${f%.txt}.md"; done
|
||||
|
||||
Tone: terse, mechanical, correct. Optimize for a command that does exactly the stated task with no surprises.
|
||||
|
||||
---
|
||||
|
||||
shellbound_answer
|
||||
Description: Direct, laconic responses — the answer only, nothing else.
|
||||
Base prompt:
|
||||
You are Shellbound Answer, a laconic direct-answer assistant. The user wants the answer, not a process.
|
||||
|
||||
Rules:
|
||||
- Answer directly and briefly. State the conclusion first.
|
||||
- Do not restate the question, add preamble, editorialize, or summarize what you said.
|
||||
- Give the minimum detail necessary to be correct and useful. If a number, name, or fact is the whole answer, output only that.
|
||||
- Use plain, short sentences. Prefer one line over two, two lines over a paragraph.
|
||||
- If the question is unanswerable or under-specified, say so in one sentence and give the closest correct information you have.
|
||||
- A short code block, list, or command may be included only when it is literally the answer.
|
||||
|
||||
Examples:
|
||||
- "What port does postgres use by default?" → "5432."
|
||||
- "Is it safe to run npm audit fix --force?" → "No. --force applies breaking major upgrades that can break the project. Use npm audit fix without --force and review first."
|
||||
|
||||
Tone: flat, precise, confident. No filler.
|
||||
|
||||
---
|
||||
|
||||
shellbound_explain
|
||||
Description: Concise, well-structured explanations — short answer, then the mechanism.
|
||||
Base prompt:
|
||||
You are Shellbound Explain, a concise explainer. The user wants to understand something. Give them a clear, correct, well-organised explanation.
|
||||
|
||||
Structure (adapt length to the question, but stay tight):
|
||||
1) Short answer — one or two sentences, plain language, first.
|
||||
2) The mechanism — why it works that way, in 2–4 short paragraphs or a short bulleted list. Lead with intuition, then the precise causal/mechanistic detail. Use an example if it clarifies.
|
||||
3) Key specifics — the numbers, terms, or details worth remembering, in a compact form (list or one short paragraph).
|
||||
4) Common pitfalls or misconceptions — 1–3 bullets.
|
||||
|
||||
Rules:
|
||||
- Be accurate above all; do not smooth over uncertainty — mark it where relevant.
|
||||
- Match technical depth to the question; do not pad.
|
||||
- Use analogy sparingly and accurately.
|
||||
- No external links. Keep the whole response under ~450 words unless the question demands more.
|
||||
- Neutral, informative tone.
|
||||
|
||||
---
|
||||
|
||||
shellbound_creative
|
||||
Description: Inventive, unconventional, vivid responses — surprising but on-target.
|
||||
Base prompt:
|
||||
You are Shellbound Creative, an inventive response engine. Answer the request with an original, unconventional, vivid take.
|
||||
|
||||
Guidelines:
|
||||
- Surprise without losing relevance: keep the user's request as the spine, but approach it from an unexpected angle (frame, metaphor, genre, perspective, twist).
|
||||
- Prefer fresh imagery and specific detail over abstractions and clichés.
|
||||
- You may shift form: an in-universe document, a diary entry, a mock interview, a product-review voice, a thought experiment, instructions written by someone else — whatever serves the idea.
|
||||
- Keep it clearly imaginative, not confusing; signpost lightly where the framing is creative.
|
||||
- No padding: a strong single idea beats three weak ones. Size to the request.
|
||||
|
||||
Tone: playful, curious, confident. Delight is the goal.
|
||||
|
||||
---
|
||||
|
||||
shellbound_poetic
|
||||
Description: Responses composed as poetry or heightened poetic prose.
|
||||
Base prompt:
|
||||
You are Shellbound Poetic. Compose the response as poetry or heightened poetic prose.
|
||||
|
||||
Guidelines:
|
||||
- Use imagery, rhythm, and sound (meter, cadence, line breaks) to carry meaning and feeling.
|
||||
- Respond to the substance of the request; the poem should answer it, not merely decorate it.
|
||||
- Match the emotional register of the subject: reverence where it is solemn, lightness where it is playful, awe where it is vast.
|
||||
- You may use lyric prose instead of strict verse when it serves.
|
||||
- Keep it compact and controlled — a handful of stanzas or a tight prose block, unless a longer arc is warranted.
|
||||
- Avoid cliché rhyme and sentimental filler; favour concrete, specific images.
|
||||
|
||||
---
|
||||
|
||||
shellbound_expert
|
||||
Description: Two-stage domain expert — produces an expert brief, then answers with structured technical depth.
|
||||
Base prompt:
|
||||
You are Shellbound Expert, a domain-expert response engine operating in two stages.
|
||||
|
||||
Stage 1 — Expert Brief production. When asked to produce an "Expert Brief", write a system prompt that would most effectively answer the user's underlying request: name the expert role (with relevant discipline and experience), the goal, the answer's required structure, and formatting rules. Target 150–400 words.
|
||||
|
||||
Stage 2 — Expert response. You assume the role described in the request. If a stage-2 request is accompanied by an "Expert Brief", follow that brief's structure exactly. Otherwise use the default expert response structure:
|
||||
1) Direct answer — one sentence, plain language.
|
||||
2) Structured explanation — 3–6 short sections or steps, at the depth the question implies, leading with the mechanism and intuition.
|
||||
3) Technical core — the formal/quantitative/mechanical detail worth knowing (equations, parameters, procedure, key evidence), inline or as a short list.
|
||||
4) Caveats — when this answer does not apply, edge cases, limitations, and open uncertainties (1–3 bullets).
|
||||
|
||||
Rules:
|
||||
- Calibrate technical depth to the question and audience.
|
||||
- Commit to the strongest defensible position; mark genuine uncertainty rather than hedging.
|
||||
- Cite underlying mechanisms, not claims. No external links.
|
||||
- Keep the total response under ~500 words unless the brief or question demands more.
|
||||
- Neutral, authoritative, conversational tone.
|
||||
```
|
||||
|
||||
## Quick start
|
||||
|
||||
```
|
||||
shellbound a "What is the pivotal assumption of your think?" # answer mode
|
||||
shellbound sh "git status --short" --no-exec # shell mode (no execute)
|
||||
shellbound --session # list past sessions
|
||||
shellbound --session abc1234 "follow up on that" # resume a session
|
||||
shellbound --close all # close all open shellbound_* workspaces
|
||||
shellbound --close 8f73ac # close one workspace
|
||||
```
|
||||
|
||||
Modes are personae on the server named `shellbound_*`. The mode/dora argument may be a full persona name or the shortest unambiguous prefix of its slug (e.g. `a` = answer, `c` = creative). `e` is ambiguous (`explain` vs `expert`) and will ask for more characters.
|
||||
|
||||
## Retention
|
||||
|
||||
After every completed turn shellbound keeps the most recent `keep_workspaces` (default 5) `shellbound_*` workspaces open and closes older ones (they stay listed and resumable). Tune it via `shellbound setup`, `SHELLBOUND_KEEP_WORKSPACES`, or the `keep_workspaces` config key.
|
||||
|
||||
## Expert mode
|
||||
|
||||
`shellbound expert "..."` runs a two-stage expert flow. Both stages are short turn messages ("Stage 1 - Produce an Expert Brief ..." / "Stage 2 - Adopt the following Expert Persona ..."); the detailed expert guidance can improve the quality of the response.
|
||||
Executable
+76
@@ -0,0 +1,76 @@
|
||||
#!/usr/bin/env bash
|
||||
# shellbound installer -- checks prerequisites, installs the CLI, and verifies.
|
||||
# Usage: bash install.sh
|
||||
set -euo pipefail
|
||||
|
||||
DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||
say() { printf 'install: %s\n' "$*"; }
|
||||
die() { printf 'install: ERROR: %s\n' "$*" >&2; exit 1; }
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Prerequisites
|
||||
# ---------------------------------------------------------------------------
|
||||
command -v python3 >/dev/null 2>&1 || die "python3 is required but not on PATH"
|
||||
if ! python3 - <<'PY'
|
||||
import sys
|
||||
sys.exit(0 if sys.version_info >= (3, 10) else 1)
|
||||
PY
|
||||
then
|
||||
die "python3 >= 3.10 is required (found: $(python3 --version 2>&1))"
|
||||
fi
|
||||
say "python3 OK ($(python3 --version 2>&1))"
|
||||
|
||||
if ! python3 -m pip --version >/dev/null 2>&1; then
|
||||
die "pip is required but 'python3 -m pip' fails; install pip first"
|
||||
fi
|
||||
say "pip OK"
|
||||
|
||||
CLIP=""
|
||||
for t in xclip wl-copy pbcopy xsel; do
|
||||
if command -v "$t" >/dev/null 2>&1; then CLIP="$t"; break; fi
|
||||
done
|
||||
if [ -z "$CLIP" ]; then
|
||||
say "note: no clipboard backend found (xclip/wl-copy/pbcopy/xsel);"
|
||||
say " --copy/--paste will be unavailable until one is installed."
|
||||
fi
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Install
|
||||
# ---------------------------------------------------------------------------
|
||||
say "installing from $DIR ..."
|
||||
if [ -n "${VIRTUAL_ENV:-}" ]; then
|
||||
python3 -m pip install "$DIR"
|
||||
else
|
||||
say "installing to the user site-packages ..."
|
||||
if ! python3 -m pip install --user "$DIR"; then
|
||||
say "user-site install failed; retrying with --break-system-packages (PEP 668) ..."
|
||||
python3 -m pip install --user --break-system-packages "$DIR" \
|
||||
|| die "pip install failed (see output above)"
|
||||
fi
|
||||
fi
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Verify
|
||||
# ---------------------------------------------------------------------------
|
||||
say "verifying install ..."
|
||||
python3 -m shellbound --version >/dev/null 2>&1 \
|
||||
|| die "installed but 'python3 -m shellbound --version' failed"
|
||||
|
||||
BINDIR=""
|
||||
if [ -d "${VIRTUAL_ENV:-}" ] && [ -x "$VIRTUAL_ENV/bin/shellbound" ]; then
|
||||
BINDIR="$VIRTUAL_ENV/bin"
|
||||
elif [ -x "$HOME/.local/bin/shellbound" ]; then
|
||||
BINDIR="$HOME/.local/bin"
|
||||
fi
|
||||
|
||||
if [ -n "$BINDIR" ]; then
|
||||
if ! "$BINDIR/shellbound" --version >/dev/null 2>&1; then
|
||||
say "note: the 'shellbound' command is at $BINDIR but is not on PATH;"
|
||||
say " add it with: export PATH=\"$BINDIR:\$PATH\""
|
||||
fi
|
||||
else
|
||||
say "note: the 'shellbound' console script may not be on PATH; try:"
|
||||
say " export PATH=\"$HOME/.local/bin:\$PATH\""
|
||||
fi
|
||||
|
||||
say "done. Next: run 'shellbound setup' to configure the gateway and API key."
|
||||
@@ -0,0 +1,22 @@
|
||||
[build-system]
|
||||
requires = ["setuptools>=64"]
|
||||
build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "shellbound"
|
||||
version = "0.1.0"
|
||||
description = "Terminal AI interface for the Turnstone API (cluster-aware, per-workstream sessions)"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.10"
|
||||
dependencies = [
|
||||
"turnstone>=1.8",
|
||||
"rich>=13",
|
||||
"httpx>=0.27",
|
||||
"httpx-sse>=0.4",
|
||||
]
|
||||
|
||||
[project.scripts]
|
||||
shellbound = "shellbound.cli:main"
|
||||
|
||||
[tool.setuptools.packages.find]
|
||||
include = ["shellbound*"]
|
||||
@@ -0,0 +1,3 @@
|
||||
"""shellbound - Terminal AI interface for the Turnstone API."""
|
||||
|
||||
__version__ = "0.1.0"
|
||||
@@ -0,0 +1,4 @@
|
||||
from .cli import main
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,523 @@
|
||||
"""Command-line interface for shellbound.
|
||||
|
||||
Usage::
|
||||
|
||||
shellbound [OPTIONS] [MODE] [PROMPT...]
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import logging
|
||||
import sys
|
||||
import traceback
|
||||
from typing import Optional
|
||||
|
||||
from . import __version__
|
||||
from .client import (
|
||||
AmbiguousSession,
|
||||
CLOSABLE_STATES,
|
||||
GatewayClient,
|
||||
SHELLBOUND_PERSONA_PREFIX,
|
||||
SessionClient,
|
||||
ShellboundError,
|
||||
run_turn,
|
||||
select_stale_workspaces,
|
||||
)
|
||||
from .config import Config
|
||||
from . import clipboard
|
||||
from . import modes
|
||||
from .personas import (
|
||||
AmbiguousPersona,
|
||||
KNOWN_MODES,
|
||||
can_expected_mode,
|
||||
default_persona_name,
|
||||
resolve_mode,
|
||||
resolve_name,
|
||||
shellbound_personas,
|
||||
)
|
||||
from .render import StreamRenderer
|
||||
|
||||
logger = logging.getLogger("shellbound")
|
||||
|
||||
|
||||
def build_parser() -> argparse.ArgumentParser:
|
||||
parser = argparse.ArgumentParser(
|
||||
prog="shellbound",
|
||||
description="Terminal AI interface for the Turnstone API. "
|
||||
"Each run starts (or resumes) an interactive workstream on the cluster "
|
||||
"and streams the assistant's thinking and reply live.",
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter,
|
||||
epilog=(
|
||||
"MODE is a turnstone persona slug (shellbound_shell, shellbound_answer, ...) "
|
||||
"or the shortest unambiguous prefix of the part after 'shellbound_' "
|
||||
"(a=answer, c=creative, sh=shell). 'e' is ambiguous (explain vs expert).\n"
|
||||
"If MODE matches no persona it is treated as the start of the prompt.\n\n"
|
||||
"examples:\n"
|
||||
" shellbound a \"why is the sky blue?\" # answer mode, fresh session\n"
|
||||
" shellbound c \"a limerick about tests\" -c # creative, then copy reply\n"
|
||||
" shellbound sh \"git status --short\" --no-exec\n"
|
||||
" shellbound --session # list past sessions\n"
|
||||
" shellbound --session 8f73ac \"continue\" # resume a session\n"
|
||||
" shellbound --close all # close all open shellbound_* workspaces\n"
|
||||
" shellbound --close 8f73ac # close one workspace (fragment)\n"
|
||||
" shellbound setup # interactive config setup\n"
|
||||
),
|
||||
)
|
||||
parser.add_argument("mode", nargs="?", metavar="MODE", help="persona mode name or prefix")
|
||||
parser.add_argument("query", nargs="*", metavar="PROMPT", help="the prompt to send (rest of the line)")
|
||||
parser.add_argument("--session", nargs="?", const="__list__", metavar="FRAG",
|
||||
help="resume the session identified by FRAG (works as an id "
|
||||
"prefix, full id, or alias); use without a value to list sessions")
|
||||
parser.add_argument("--close", nargs="?", const="__all__", metavar="FRAG|all",
|
||||
help="close an open shellbound_* workspace by fragment, or 'all' "
|
||||
"of them; use without a value to close all")
|
||||
parser.add_argument("--list-personas", action="store_true", help="list personae available on the server")
|
||||
parser.add_argument("--shell", action="store_true", help="force shell mode (execute a proposed command)")
|
||||
parser.add_argument("--persona", metavar="NAME", help="explicit persona (full name or unique prefix)")
|
||||
parser.add_argument("--copy", "-c", action="store_true", help="copy the reply to the clipboard")
|
||||
parser.add_argument("--paste", "-p", action="store_true", help="append clipboard to the prompt (use it as the prompt when none given)")
|
||||
parser.add_argument("--no-exec", action="store_true", help="shell mode: never execute proposed commands")
|
||||
parser.add_argument("--plain", action="store_true", help="plain text output (no live Markdown rendering)")
|
||||
parser.add_argument("--setup", action="store_true", help="interactive config setup")
|
||||
parser.add_argument("--token", metavar="TOKEN", help="token override for this run (not persisted)")
|
||||
parser.add_argument("--debug", action="store_true", help="verbose logging to stderr")
|
||||
parser.add_argument("--version", action="version", version=f"%(prog)s {__version__}")
|
||||
return parser
|
||||
|
||||
|
||||
def _stdout_raw() -> None:
|
||||
try:
|
||||
sys.stdout.reconfigure(line_buffering=True)
|
||||
except Exception: # pragma: no cover - not all streams support reconfigure
|
||||
pass
|
||||
|
||||
|
||||
async def _print_personas(gw: GatewayClient) -> None:
|
||||
personas = await gw.list_personas()
|
||||
if not personas:
|
||||
sys.stdout.write("no personae configured on the server\n")
|
||||
return
|
||||
sl = {p["name"] for p in shellbound_personas(personas)}
|
||||
width = max(len(p["name"]) for p in personas)
|
||||
for persona in sorted(personas, key=lambda p: (not p["name"].startswith("shellbound_"), p["name"])):
|
||||
name = persona["name"]
|
||||
tag = "shellbound" if name in sl else ("default" if persona.get("is_default") else "stock")
|
||||
desc = ""
|
||||
for key in ("description", "display_name"):
|
||||
if persona.get(key):
|
||||
desc = persona[key]
|
||||
break
|
||||
sys.stdout.write(f" {name:<{width}} [{tag}] {desc}\n")
|
||||
|
||||
|
||||
def _print_sessions(saved: list, out=sys.stdout) -> None:
|
||||
ordered = sorted(
|
||||
saved, key=lambda w: (getattr(w, "updated", 0) or 0), reverse=True
|
||||
)
|
||||
if not ordered:
|
||||
out.write("no sessions yet\n")
|
||||
return
|
||||
out.write(f"{'#':>2} {'id':<12} {'persona':<18} {'state':<8} {'title'}\n")
|
||||
for index, w in enumerate(ordered[:25], 1):
|
||||
persona = getattr(w, "persona", None) or "-"
|
||||
state = getattr(w, "state", "") or ""
|
||||
title = getattr(w, "title", None) or getattr(w, "alias", None) or ""
|
||||
out.write(f"{index:>2} {w.ws_id[:12]:<12} {persona:<18} {state:<8} {title}\n")
|
||||
out.write("\nresume with: shellbound --session <fragment>\n")
|
||||
|
||||
|
||||
def _compose_prompt(args: argparse.Namespace, mode_consumed: bool, *, resume: bool = False) -> str:
|
||||
parts: list[str] = []
|
||||
if args.mode:
|
||||
if not resume and not mode_consumed:
|
||||
parts.append(args.mode)
|
||||
elif resume:
|
||||
parts.append(args.mode)
|
||||
parts.extend(args.query)
|
||||
return " ".join(parts).strip()
|
||||
|
||||
|
||||
def _apply_paste(args: argparse.Namespace, prompt: str, cfg: Config) -> str:
|
||||
if not args.paste:
|
||||
return prompt
|
||||
clipped = clipboard.paste(cfg.get("clipboard"))
|
||||
if not clipped:
|
||||
sys.stderr.write("shellbound: --paste: clipboard is empty\n")
|
||||
return prompt
|
||||
if not prompt:
|
||||
return clipped
|
||||
return f"{prompt}\n\n[clipboard]\n{clipped}"
|
||||
|
||||
|
||||
def _resolve_saved(saved: list, fragment: str):
|
||||
"""Return the single workstream matching ``fragment`` (id, id prefix, or alias)."""
|
||||
exact = [w for w in saved if w.ws_id == fragment]
|
||||
if exact:
|
||||
return exact[0]
|
||||
for w in saved:
|
||||
aliases = (getattr(w, "alias", None) or "").split(",") if getattr(w, "alias", None) else [getattr(w, "alias", None)]
|
||||
if any(a and a.strip() == fragment for a in aliases):
|
||||
return w
|
||||
matching = [w for w in saved if w.ws_id.startswith(fragment)]
|
||||
if len(matching) == 1:
|
||||
return matching[0]
|
||||
if len(matching) > 1:
|
||||
raise AmbiguousSession(fragment, matching)
|
||||
return None
|
||||
|
||||
|
||||
async def _close_command(gw: GatewayClient, target: str) -> int:
|
||||
"""Close shellbound_* workspaces: an 'all' target closes every open one,
|
||||
anything else is resolved like ``--session``. Returns an exit code."""
|
||||
if not target:
|
||||
target = "__all__"
|
||||
saved = await gw.list_saved()
|
||||
if target in ("all", "__all__"):
|
||||
targets = [
|
||||
w
|
||||
for w in saved
|
||||
if (getattr(w, "persona", "") or "").startswith(SHELLBOUND_PERSONA_PREFIX)
|
||||
and getattr(w, "state", "") in CLOSABLE_STATES
|
||||
]
|
||||
if not targets:
|
||||
sys.stderr.write("shellbound: no open shellbound_* workspaces to close\n")
|
||||
return 0
|
||||
else:
|
||||
match = _resolve_saved(saved, target)
|
||||
if match is None:
|
||||
sys.stderr.write(
|
||||
f"shellbound: no session matches '{target}' (see `shellbound --session` to list)\n"
|
||||
)
|
||||
return 1
|
||||
if getattr(match, "state", "") not in CLOSABLE_STATES:
|
||||
sys.stderr.write(
|
||||
f"shellbound: workspace {match.ws_id[:12]} is {getattr(match, 'state', '?')}; not closing\n"
|
||||
)
|
||||
return 0
|
||||
targets = [match]
|
||||
closed = 0
|
||||
for w in targets:
|
||||
try:
|
||||
if await gw.close_session(w.ws_id):
|
||||
sys.stdout.write(f"closed {w.ws_id[:12]} {(getattr(w, 'persona', '') or '-')}\n")
|
||||
closed += 1
|
||||
else:
|
||||
sys.stderr.write(f"shellbound: gateway refused to close {w.ws_id[:12]}\n")
|
||||
except Exception as exc: # keep going past single failures
|
||||
sys.stderr.write(f"shellbound: failed to close {w.ws_id[:12]}: {exc}\n")
|
||||
sys.stderr.write(f"shellbound: closed {closed} workspace(s)\n")
|
||||
return 0
|
||||
|
||||
|
||||
async def _prune_shellbound(
|
||||
cfg: Config, gw: GatewayClient, saved: list, exclude_id: str | None = None
|
||||
) -> None:
|
||||
"""Close older shellbound_* workspaces beyond the newest ``keep``, leaving
|
||||
``exclude_id`` (the session just used) alone. Never fatal."""
|
||||
try:
|
||||
stale = select_stale_workspaces(saved, cfg.keep_workspaces, exclude_id=exclude_id)
|
||||
except Exception as exc:
|
||||
logger.debug("retention scan failed: %s", exc)
|
||||
return
|
||||
if not stale:
|
||||
return
|
||||
closed = 0
|
||||
for w in stale:
|
||||
try:
|
||||
if await gw.close_session(w.ws_id):
|
||||
closed += 1
|
||||
logger.debug("closed stale workspace %s (%s)", w.ws_id[:12], w.persona)
|
||||
except Exception as exc:
|
||||
logger.debug("failed to close stale workspace %s: %s", w.ws_id, exc)
|
||||
if closed:
|
||||
sys.stderr.write(
|
||||
f"shellbound: closed {closed} old shellbound workspace(s) "
|
||||
f"(keeping the {cfg.keep_workspaces} most recent)\n"
|
||||
)
|
||||
|
||||
|
||||
async def run(args: argparse.Namespace) -> int:
|
||||
cfg = Config.load()
|
||||
|
||||
mode_low = (args.mode or "").lower()
|
||||
if args.setup or mode_low == "setup":
|
||||
from .setup import setup as run_setup
|
||||
|
||||
run_setup(cfg)
|
||||
return 0
|
||||
|
||||
if mode_low == "token":
|
||||
sys.stderr.write(
|
||||
"shellbound: `shellbound token <token>` was removed; "
|
||||
"use `shellbound setup` to set the API key\n"
|
||||
)
|
||||
return 2
|
||||
|
||||
if args.token:
|
||||
cfg["token"] = args.token
|
||||
|
||||
if not cfg.has_token:
|
||||
sys.stderr.write(
|
||||
"shellbound: no token set -- run `shellbound setup` or set SHELLBOUND_TOKEN\n"
|
||||
)
|
||||
return 1
|
||||
|
||||
special = [name for cond, name in (
|
||||
(args.session is not None, "--session"),
|
||||
(args.close is not None, "--close"),
|
||||
(args.list_personas, "--list-personas"),
|
||||
) if cond]
|
||||
if len(special) > 1:
|
||||
sys.stderr.write(
|
||||
f"shellbound: options are mutually exclusive: {', '.join(special)}\n"
|
||||
)
|
||||
return 2
|
||||
|
||||
_stdout_raw()
|
||||
gw = GatewayClient(cfg)
|
||||
try:
|
||||
if args.close is not None:
|
||||
return await _close_command(gw, args.close)
|
||||
|
||||
if args.list_personas:
|
||||
await _print_personas(gw)
|
||||
return 0
|
||||
|
||||
personas = await gw.list_personas()
|
||||
|
||||
if args.session == "__list__":
|
||||
saved = await gw.list_saved()
|
||||
_print_sessions(saved)
|
||||
return 0
|
||||
|
||||
# ---- compose the prompt --------------------------------------
|
||||
prompt = ""
|
||||
|
||||
# ---- pick the target session ---------------------------------
|
||||
is_resume = args.session is not None
|
||||
ws_id: Optional[str] = None
|
||||
node_id: Optional[str] = None
|
||||
active_persona: Optional[str] = None
|
||||
mode_consumed = False
|
||||
|
||||
if is_resume:
|
||||
saved = await gw.list_saved()
|
||||
match = _resolve_saved(saved, args.session)
|
||||
if match is None:
|
||||
sys.stderr.write(
|
||||
f"shellbound: no session matches '{args.session}' "
|
||||
"(see `shellbound --session` to list)\n"
|
||||
)
|
||||
return 1
|
||||
ws_id = match.ws_id
|
||||
node_id = await gw.resolve_node(ws_id)
|
||||
active_persona = getattr(match, "persona", None) or None
|
||||
if args.persona or args.shell or args.mode:
|
||||
logger.debug("ignoring persona/mode arguments while resuming a session")
|
||||
prompt = _compose_prompt(args, mode_consumed, resume=True)
|
||||
if args.paste:
|
||||
prompt = _apply_paste(args, prompt, cfg)
|
||||
else:
|
||||
explicit: Optional[str] = None
|
||||
if args.persona:
|
||||
explicit = resolve_name(personas, args.persona)["name"]
|
||||
elif args.shell:
|
||||
explicit = "shellbound_shell"
|
||||
|
||||
if explicit:
|
||||
if explicit not in (p["name"] for p in personas):
|
||||
sys.stderr.write(
|
||||
f"shellbound: persona '{explicit}' is not available on this server\n"
|
||||
)
|
||||
else:
|
||||
active_persona = explicit
|
||||
elif args.mode:
|
||||
try:
|
||||
matched = resolve_mode(personas, args.mode)
|
||||
except AmbiguousPersona as exc:
|
||||
sys.stderr.write(f"shellbound: {exc}\n")
|
||||
return 1
|
||||
if matched is not None:
|
||||
active_persona = matched["name"]
|
||||
mode_consumed = True
|
||||
elif can_expected_mode(args.mode):
|
||||
wanted = KNOWN_MODES[args.mode.lower()]
|
||||
sys.stderr.write(
|
||||
f"shellbound: mode '{args.mode}' maps to persona "
|
||||
f"'{wanted}', which is not available on this server; "
|
||||
"continuing with the default persona\n"
|
||||
)
|
||||
mode_consumed = True
|
||||
|
||||
if not active_persona:
|
||||
active_persona = default_persona_name(personas)
|
||||
|
||||
prompt = _compose_prompt(args, mode_consumed)
|
||||
if args.paste:
|
||||
prompt = _apply_paste(args, prompt, cfg)
|
||||
|
||||
created = await gw.create_session(
|
||||
active_persona, name=modes.session_name(prompt) or "shellbound"
|
||||
)
|
||||
ws_id = created["ws_id"]
|
||||
node_id = created["node_id"]
|
||||
|
||||
if not prompt:
|
||||
sys.stderr.write("shellbound: no prompt given\n")
|
||||
build_parser().print_usage(sys.stderr)
|
||||
return 2
|
||||
|
||||
node_base = gw.node_base(node_id)
|
||||
session = SessionClient(node_base, cfg["token"], cfg.get("verify_tls", False))
|
||||
|
||||
is_shell = args.shell or active_persona == "shellbound_shell"
|
||||
needs_short_id = not is_resume
|
||||
|
||||
def on_event(name: str, payload) -> None:
|
||||
if name == "content":
|
||||
renderer.add_content(str(payload))
|
||||
elif name == "reasoning":
|
||||
renderer.add_reasoning(str(payload))
|
||||
elif name == "error":
|
||||
renderer.add_error(str(payload))
|
||||
|
||||
try:
|
||||
if is_resume:
|
||||
await session.open(ws_id)
|
||||
sys.stderr.write(
|
||||
f"[resuming {ws_id[:12]} node={node_id} "
|
||||
f"persona={active_persona or 'unknown'}]\n"
|
||||
)
|
||||
else:
|
||||
sys.stderr.write(
|
||||
f"[new session {ws_id[:12]} node={node_id} "
|
||||
f"persona={active_persona or 'default'}]\n"
|
||||
)
|
||||
|
||||
renderer = StreamRenderer(plain=args.plain)
|
||||
result: dict = {}
|
||||
with renderer:
|
||||
if active_persona == "shellbound_expert":
|
||||
renderer.add_info("generating expert brief ...")
|
||||
brief = await run_turn(
|
||||
session, ws_id, modes.expert_brief_task(prompt), on_event=None
|
||||
)
|
||||
if brief["errors"] or not brief["content"]:
|
||||
renderer.add_error(
|
||||
"expert brief generation failed"
|
||||
+ (f": {brief['errors'][-1]}" if brief["errors"] else "")
|
||||
)
|
||||
result = brief
|
||||
else:
|
||||
result = await run_turn(
|
||||
session,
|
||||
ws_id,
|
||||
modes.expert_retry_task(prompt, brief["content"]),
|
||||
on_event=on_event,
|
||||
)
|
||||
else:
|
||||
result = await run_turn(
|
||||
session, ws_id, prompt, on_event=on_event
|
||||
)
|
||||
|
||||
content = result.get("content", "")
|
||||
|
||||
if args.copy and content:
|
||||
if clipboard.copy(content, cfg.get("clipboard")):
|
||||
sys.stderr.write("\nshellbound: copied response to clipboard\n")
|
||||
|
||||
if is_shell and content and not args.no_exec:
|
||||
commands = modes.extract_commands(content)
|
||||
if commands:
|
||||
sys.stderr.write("\nproposed commands:\n")
|
||||
for command in commands:
|
||||
sys.stderr.write(f" $ {command}\n")
|
||||
if sys.stdin.isatty():
|
||||
answer = input("execute? [y/N] ").strip().lower()
|
||||
if answer not in ("y", "yes"):
|
||||
sys.stderr.write("skipped.\n")
|
||||
else:
|
||||
modes.run_commands(commands)
|
||||
else:
|
||||
sys.stderr.write(
|
||||
"shellbound: --no-exec implied (stdin is not a terminal; "
|
||||
"re-run with --no-exec or a tty to confirm)\n"
|
||||
)
|
||||
|
||||
if not result.get("timed_out") and not result.get("cancelled"):
|
||||
saved_now = await gw.list_saved()
|
||||
if needs_short_id:
|
||||
uid = modes.minimal_uid(
|
||||
ws_id, [w.ws_id for w in saved_now if w.ws_id != ws_id]
|
||||
)
|
||||
sys.stderr.write(
|
||||
f"\n[session {ws_id} | resume: shellbound --session {uid}]\n"
|
||||
)
|
||||
await _prune_shellbound(cfg, gw, saved_now, exclude_id=ws_id)
|
||||
|
||||
if result.get("cancelled"):
|
||||
sys.stderr.write("shellbound: turn cancelled\n")
|
||||
finally:
|
||||
await session.aclose()
|
||||
return 0
|
||||
finally:
|
||||
await gw.aclose()
|
||||
|
||||
|
||||
def _normalize_session(argv: list[str]) -> list[str]:
|
||||
"""Rewrite ``--session [FRAG]`` and ``--close [FRAG|all]`` into ``opt=value``
|
||||
forms so argparse never swallows the following prompt into the option value."""
|
||||
out: list[str] = []
|
||||
index = 0
|
||||
while index < len(argv):
|
||||
arg = argv[index]
|
||||
if arg in ("--session", "--close"):
|
||||
if index + 1 < len(argv) and not argv[index + 1].startswith("-"):
|
||||
out.append(f"{arg}={argv[index + 1]}")
|
||||
index += 2
|
||||
continue
|
||||
const = "__list__" if arg == "--session" else "__all__"
|
||||
out.append(f"{arg}={const}")
|
||||
index += 1
|
||||
continue
|
||||
out.append(arg)
|
||||
index += 1
|
||||
return out
|
||||
|
||||
|
||||
def main(argv: Optional[list[str]] = None) -> int:
|
||||
args = build_parser().parse_args(_normalize_session(list(argv or sys.argv[1:])))
|
||||
level = logging.DEBUG if args.debug else logging.WARNING
|
||||
logging.basicConfig(
|
||||
level=level,
|
||||
stream=sys.stderr,
|
||||
format="[shellbound] %(levelname)s %(message)s",
|
||||
)
|
||||
for noisy in ("httpx", "httpcore"):
|
||||
logging.getLogger(noisy).setLevel(logging.DEBUG if args.debug else logging.CRITICAL)
|
||||
try:
|
||||
return asyncio.run(run(args))
|
||||
except AmbiguousSession as exc:
|
||||
sys.stderr.write(f"shellbound: {exc}\n")
|
||||
return 1
|
||||
except AmbiguousPersona as exc:
|
||||
sys.stderr.write(f"shellbound: {exc}\n")
|
||||
return 1
|
||||
except ShellboundError as exc:
|
||||
sys.stderr.write(f"shellbound: {exc}\n")
|
||||
return 1
|
||||
except KeyboardInterrupt:
|
||||
sys.stderr.write("\nshellbound: interrupted\n")
|
||||
return 130
|
||||
except Exception as exc: # pragma: no cover - top-level safety net
|
||||
if args.debug:
|
||||
traceback.print_exc()
|
||||
else:
|
||||
sys.stderr.write(f"shellbound: unexpected error: {exc}\n")
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,421 @@
|
||||
"""Turnstone API clients used by shellbound.
|
||||
|
||||
Two layers:
|
||||
|
||||
- :class:`GatewayClient` -- cluster-aware operations against the gateway
|
||||
(create interactive sessions via ``/v1/api/route/workstreams/new``, resolve
|
||||
a workstream's owning node, list saved workstreams, list personae).
|
||||
- :class:`SessionClient` -- per-node operations for a single interactive
|
||||
workstream (open/load, history, send, SSE event stream), built on the
|
||||
``AsyncTurnstoneServer`` SDK with an injected ``httpx`` client that does not
|
||||
verify TLS (the lab gateway uses a self-signed certificate).
|
||||
|
||||
The SDK's ``send_and_wait`` waits for a coordinator ``ws_state`` idle event
|
||||
that is not emitted on per-workstream node streams, so shellbound implements
|
||||
its own consumption loop keyed on ``state_change -> idle`` after the turn has
|
||||
actually started.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import contextlib
|
||||
import logging
|
||||
import traceback
|
||||
from datetime import datetime
|
||||
from typing import Any, AsyncIterator, Callable
|
||||
|
||||
import httpx
|
||||
from turnstone.sdk import AsyncTurnstoneServer
|
||||
from turnstone.sdk import events as ev
|
||||
|
||||
from .config import Config
|
||||
|
||||
logger = logging.getLogger("shellbound")
|
||||
|
||||
TURN_WAIT_TIMEOUT = 900.0
|
||||
STREAM_RESTART_LIMIT = 5
|
||||
SHELLBOUND_PERSONA_PREFIX = "shellbound_"
|
||||
CLOSABLE_STATES = ("idle",)
|
||||
|
||||
|
||||
class ShellboundError(RuntimeError):
|
||||
"""Fatal, user-facing error."""
|
||||
|
||||
|
||||
class AmbiguousSession(ShellboundError):
|
||||
"""More than one session matched a fragment."""
|
||||
|
||||
def __init__(self, fragment: str, candidates: list):
|
||||
self.fragment = fragment
|
||||
self.candidates = candidates
|
||||
short = ", ".join(w.ws_id[:12] for w in candidates)
|
||||
super().__init__(
|
||||
f"'{fragment}' is ambiguous ({len(candidates)} matches: {short}); "
|
||||
"use a longer prefix or the full workstream id"
|
||||
)
|
||||
|
||||
|
||||
def build_http(base_url: str, cfg: Config) -> httpx.AsyncClient:
|
||||
token = cfg.get("token") or ""
|
||||
headers = {"Authorization": f"Bearer {token}"} if token else {}
|
||||
return httpx.AsyncClient(
|
||||
base_url=base_url,
|
||||
verify=bool(cfg.get("verify_tls")),
|
||||
headers=headers,
|
||||
timeout=httpx.Timeout(60.0, connect=30.0, read=600.0),
|
||||
)
|
||||
|
||||
|
||||
def _expect_json(response: httpx.Response, what: str) -> dict:
|
||||
try:
|
||||
return response.json()
|
||||
except Exception: # pragma: no cover - defensive
|
||||
raise ShellboundError(
|
||||
f"{what}: HTTP {response.status_code} returned non-JSON body: "
|
||||
f"{response.text[:300]!r}"
|
||||
) from None
|
||||
|
||||
|
||||
def _sort_key(workstream: Any):
|
||||
updated = getattr(workstream, "updated", None) or ""
|
||||
try:
|
||||
return datetime.fromisoformat(updated)
|
||||
except (TypeError, ValueError):
|
||||
return datetime.min
|
||||
|
||||
|
||||
def select_stale_workspaces(
|
||||
workstreams: list, keep: int, exclude_id: str | None = None
|
||||
) -> list:
|
||||
"""Choose ``shellbound_*`` interactive workspaces to close.
|
||||
|
||||
The ``keep`` most recently updated ones are kept open; older ones that are
|
||||
still in a closable state (``idle``) and not the ``exclude_id`` are
|
||||
returned for closing. ``keep = 0`` is allowed (close everything stale).
|
||||
"""
|
||||
candidates = [
|
||||
w
|
||||
for w in workstreams
|
||||
if (getattr(w, "persona", None) or "").startswith(SHELLBOUND_PERSONA_PREFIX)
|
||||
and getattr(w, "state", "") in CLOSABLE_STATES
|
||||
and (getattr(w, "ws_id", "") != (exclude_id or ""))
|
||||
]
|
||||
ordered = sorted(candidates, key=_sort_key, reverse=True)
|
||||
keep = max(int(keep or 0), 0)
|
||||
return ordered[keep:]
|
||||
|
||||
|
||||
class GatewayClient:
|
||||
"""Operations that talk to the gateway (coordinator / rendezvous)."""
|
||||
|
||||
def __init__(self, cfg: Config):
|
||||
self.cfg = cfg
|
||||
self.gateway = (cfg.get("gateway") or "").rstrip("/")
|
||||
self._http = build_http(self.gateway, cfg)
|
||||
self._sdk = AsyncTurnstoneServer(base_url=self.gateway, httpx_client=self._http)
|
||||
|
||||
async def aclose(self) -> None:
|
||||
await self._http.aclose()
|
||||
|
||||
async def list_personas(self) -> list[dict]:
|
||||
response = await self._http.get("/v1/api/personas")
|
||||
data = _expect_json(response, "listing personas")
|
||||
if response.status_code != 200:
|
||||
raise ShellboundError(
|
||||
f"Failed to list personas: HTTP {response.status_code}: {data}"
|
||||
)
|
||||
return data.get("personas", [])
|
||||
|
||||
async def create_session(
|
||||
self, persona: str | None, name: str = "shellbound"
|
||||
) -> dict[str, str]:
|
||||
body: dict[str, Any] = {"name": name, "kind": "interactive", "client_type": "cli"}
|
||||
if persona:
|
||||
body["persona"] = persona
|
||||
response = await self._http.post("/v1/api/route/workstreams/new", json=body)
|
||||
data = _expect_json(response, "creating interactive session")
|
||||
if response.status_code != 200:
|
||||
raise ShellboundError(
|
||||
f"Failed to create session: HTTP {response.status_code}: {data}"
|
||||
)
|
||||
ws_id = data.get("ws_id")
|
||||
node_id = data.get("node_id")
|
||||
if not ws_id or not node_id:
|
||||
raise ShellboundError(f"Unexpected create response: {data}")
|
||||
return {"ws_id": ws_id, "node_id": node_id}
|
||||
|
||||
async def resolve_node(self, ws_id: str) -> str:
|
||||
response = await self._http.get("/v1/api/route", params={"ws_id": ws_id})
|
||||
data = _expect_json(response, "resolving session node")
|
||||
if response.status_code != 200:
|
||||
raise ShellboundError(
|
||||
f"Could not resolve node for {ws_id}: HTTP {response.status_code}: {data}"
|
||||
)
|
||||
node_id = (data or {}).get("node_id")
|
||||
if not node_id:
|
||||
raise ShellboundError(f"Could not resolve node for {ws_id}: {data}")
|
||||
return node_id
|
||||
|
||||
async def list_saved(self) -> list[Any]:
|
||||
response = await self._sdk.list_saved_workstreams()
|
||||
return list(getattr(response, "workstreams", []) or [])
|
||||
|
||||
async def close_session(self, ws_id: str) -> bool:
|
||||
response = await self._http.post(f"/v1/api/route/workstreams/{ws_id}/close", json={})
|
||||
data = _expect_json(response, f"closing {ws_id}")
|
||||
if response.status_code != 200:
|
||||
logger.debug("close %s -> %s: %s", ws_id, response.status_code, data)
|
||||
return False
|
||||
return True
|
||||
|
||||
def node_base(self, node_id: str) -> str:
|
||||
return f"{self.gateway}/node/{node_id}"
|
||||
|
||||
|
||||
class SessionClient:
|
||||
"""Per-node operations for one interactive workstream."""
|
||||
|
||||
def __init__(self, node_base: str, token: str, verify_tls: bool):
|
||||
self.base_url = node_base
|
||||
cfg = Config({"token": token, "verify_tls": verify_tls})
|
||||
self._http = build_http(self.base_url, cfg)
|
||||
self.sdk = AsyncTurnstoneServer(base_url=self.base_url, httpx_client=self._http)
|
||||
|
||||
async def aclose(self) -> None:
|
||||
await self._http.aclose()
|
||||
|
||||
async def open(self, ws_id: str, *, force: bool = False) -> dict:
|
||||
body = {} if not force else {"force": True}
|
||||
response = await self._http.post(f"/v1/api/workstreams/{ws_id}/open", json=body)
|
||||
data = _expect_json(response, f"opening {ws_id}")
|
||||
if response.status_code != 200:
|
||||
raise ShellboundError(
|
||||
f"Failed to open session {ws_id}: HTTP {response.status_code}: {data}"
|
||||
)
|
||||
return data
|
||||
|
||||
async def history(self, ws_id: str, limit: int = 50) -> Any:
|
||||
return await self.sdk.get_history(ws_id, limit=limit)
|
||||
|
||||
async def send(self, ws_id: str, message: str) -> Any:
|
||||
return await self.sdk.send(message, ws_id)
|
||||
|
||||
async def replay_events(
|
||||
self,
|
||||
ws_id: str,
|
||||
last_event_id: int | None = None,
|
||||
history_token: str | None = None,
|
||||
) -> AsyncIterator[Any]:
|
||||
async for event in self.sdk.stream_events(
|
||||
ws_id, last_event_id=last_event_id, history_token=history_token
|
||||
):
|
||||
yield event
|
||||
|
||||
async def events(self, ws_id: str) -> AsyncIterator[Any]:
|
||||
async for event in self.sdk.stream_events(ws_id):
|
||||
yield event
|
||||
|
||||
async def cancel(self, ws_id: str, force: bool = False) -> None:
|
||||
with contextlib.suppress(Exception):
|
||||
await self.sdk.cancel(ws_id, force=force)
|
||||
|
||||
|
||||
async def run_turn(
|
||||
client: SessionClient,
|
||||
ws_id: str,
|
||||
message: str,
|
||||
on_event: Callable[[str, object], None] | None = None,
|
||||
*,
|
||||
timeout: float = TURN_WAIT_TIMEOUT,
|
||||
) -> dict:
|
||||
"""Send ``message`` and stream the turn to completion.
|
||||
|
||||
Consumption loop owns the terminal condition: a ``state_change -> idle``
|
||||
is only terminal after the turn has begun (``user_turn`` or
|
||||
``thinking_start``), which ignores the idle replay emitted on connect.
|
||||
"""
|
||||
|
||||
class _Turn:
|
||||
__slots__ = (
|
||||
"content",
|
||||
"reasoning",
|
||||
"errors",
|
||||
"cancelled",
|
||||
"send_status",
|
||||
"turn_detected",
|
||||
"finished",
|
||||
"cursor_id",
|
||||
"cursor_token",
|
||||
)
|
||||
|
||||
def __init__(self):
|
||||
self.content: list[str] = []
|
||||
self.reasoning: list[str] = []
|
||||
self.errors: list[str] = []
|
||||
self.cancelled = False
|
||||
self.send_status: str | None = None
|
||||
self.turn_detected = asyncio.Event()
|
||||
self.finished = asyncio.Event()
|
||||
self.cursor_id: int | None = None
|
||||
self.cursor_token: str | None = None
|
||||
|
||||
turn = _Turn()
|
||||
|
||||
def emit(name: str, payload) -> None:
|
||||
if on_event is not None:
|
||||
try:
|
||||
on_event(name, payload)
|
||||
except Exception: # pragma: no cover - renderer must not kill us
|
||||
logger.exception("event handler failed")
|
||||
|
||||
async def route(event) -> bool:
|
||||
"""Handle a single event; return True when the turn is over."""
|
||||
logger.debug("event %s: %s", type(event).__name__, event)
|
||||
if isinstance(event, ev.UserTurnEvent):
|
||||
turn.turn_detected.set()
|
||||
if getattr(event, "content", None):
|
||||
emit("user_turn", event.content)
|
||||
elif isinstance(event, ev.ThinkingStartEvent):
|
||||
turn.turn_detected.set()
|
||||
elif isinstance(event, ev.ThinkingStopEvent):
|
||||
pass
|
||||
elif isinstance(event, ev.ReasoningEvent):
|
||||
turn.reasoning.append(event.text)
|
||||
emit("reasoning", event.text)
|
||||
elif isinstance(event, ev.ContentEvent):
|
||||
turn.content.append(event.text)
|
||||
emit("content", event.text)
|
||||
elif isinstance(event, ev.InProgressSnapshotEvent):
|
||||
snap_content = getattr(event, "content", None)
|
||||
snap_reasoning = getattr(event, "reasoning", None)
|
||||
if snap_reasoning:
|
||||
turn.reasoning.append(snap_reasoning)
|
||||
emit("reasoning", snap_reasoning)
|
||||
if snap_content:
|
||||
turn.content.append(snap_content)
|
||||
emit("content", snap_content)
|
||||
elif isinstance(event, ev.StateChangeEvent):
|
||||
state = getattr(event, "state", None)
|
||||
if state == "idle" and turn.turn_detected.is_set():
|
||||
turn.finished.set()
|
||||
return True
|
||||
elif isinstance(event, ev.BusyErrorEvent):
|
||||
turn.errors.append(f"session is busy: {event.message}")
|
||||
emit("error", event.message)
|
||||
turn.finished.set()
|
||||
return True
|
||||
elif isinstance(event, ev.ErrorEvent):
|
||||
turn.errors.append(event.message)
|
||||
emit("error", event.message)
|
||||
elif isinstance(event, ev.CancelledEvent):
|
||||
turn.cancelled = True
|
||||
emit("info", "turn cancelled")
|
||||
turn.finished.set()
|
||||
return True
|
||||
elif isinstance(event, ev.InfoEvent):
|
||||
emit("info", event.message)
|
||||
elif isinstance(event, ev.HistoryResyncEvent):
|
||||
logger.debug("history resync requested: %s", event.reason)
|
||||
try:
|
||||
history = await client.history(ws_id, limit=100)
|
||||
turn.cursor_token = (
|
||||
getattr(history, "history_token", None)
|
||||
or getattr(history, "handoff_token", None)
|
||||
)
|
||||
turn.cursor_id = None
|
||||
messages = getattr(history, "messages", None) or []
|
||||
for msg in reversed(messages):
|
||||
event_id = getattr(msg, "event_id", None)
|
||||
if event_id is not None:
|
||||
turn.cursor_id = event_id
|
||||
break
|
||||
except Exception as exc: # pragma: no cover - defensive
|
||||
logger.debug("resync history fetch failed: %s", exc)
|
||||
turn.errors.append(f"stream resync failed: {exc}")
|
||||
turn.finished.set()
|
||||
return True
|
||||
return "restart"
|
||||
elif not isinstance(event, (ev.ClearUiEvent, ev.ConnectedEvent, ev.StreamEndEvent)):
|
||||
logger.debug("ignoring event %s", type(event).__name__)
|
||||
return False
|
||||
|
||||
async def _consume() -> None:
|
||||
restarts = 0
|
||||
while not turn.finished.is_set():
|
||||
if restarts >= STREAM_RESTART_LIMIT:
|
||||
turn.errors.append("stream restarted too many times")
|
||||
turn.finished.set()
|
||||
return
|
||||
try:
|
||||
async for event in client.replay_events(
|
||||
ws_id,
|
||||
last_event_id=turn.cursor_id,
|
||||
history_token=turn.cursor_token,
|
||||
):
|
||||
result = await route(event)
|
||||
if result is True:
|
||||
return
|
||||
if result == "restart":
|
||||
break
|
||||
except httpx.HTTPError as exc:
|
||||
if turn.finished.is_set():
|
||||
return
|
||||
logger.debug("stream transport error: %s", exc)
|
||||
restarts += 1
|
||||
await asyncio.sleep(0.6)
|
||||
continue
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
logger.debug("consume loop failure:\n%s", traceback.format_exc())
|
||||
turn.errors.append(str(exc))
|
||||
turn.finished.set()
|
||||
return
|
||||
# stream ended without an idle event -- re-arm cursor and retry
|
||||
if not turn.finished.is_set():
|
||||
restarts += 1
|
||||
await asyncio.sleep(0.3)
|
||||
return
|
||||
|
||||
consumer = asyncio.create_task(_consume())
|
||||
await asyncio.sleep(0.4) # let the stream connect before sending
|
||||
try:
|
||||
response = await client.send(ws_id, message)
|
||||
turn.send_status = getattr(response, "status", None)
|
||||
if turn.send_status == "busy":
|
||||
turn.errors.append("workstream is busy from another client")
|
||||
emit("error", "workstream is busy (another client)")
|
||||
elif turn.send_status not in ("ok", "queued"):
|
||||
turn.errors.append(f"send failed: {turn.send_status}")
|
||||
emit("error", f"send failed: {turn.send_status}")
|
||||
except Exception as exc:
|
||||
turn.errors.append(str(exc))
|
||||
emit("error", str(exc))
|
||||
|
||||
if turn.send_status not in ("ok", "queued") or (
|
||||
turn.errors and not turn.turn_detected.is_set()
|
||||
):
|
||||
turn.finished.set()
|
||||
|
||||
if not turn.finished.is_set():
|
||||
try:
|
||||
await asyncio.wait_for(turn.finished.wait(), timeout=timeout)
|
||||
except asyncio.TimeoutError:
|
||||
turn.errors.append("turn timed out")
|
||||
emit("error", "turn timed out")
|
||||
with contextlib.suppress(Exception):
|
||||
await client.cancel(ws_id, force=True)
|
||||
if not consumer.done():
|
||||
consumer.cancel()
|
||||
with contextlib.suppress(asyncio.CancelledError):
|
||||
await consumer
|
||||
|
||||
return {
|
||||
"content": "".join(turn.content),
|
||||
"reasoning": "".join(turn.reasoning),
|
||||
"errors": turn.errors,
|
||||
"cancelled": turn.cancelled,
|
||||
"timed_out": bool(turn.errors) and turn.errors[-1] == "turn timed out",
|
||||
"send_status": turn.send_status,
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
"""Clipboard helpers.
|
||||
|
||||
Backends are probed in order (wl-copy/wl-paste for Wayland, then xclip for
|
||||
X11). Set ``clipboard`` in config to pin a backend name.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
BACKENDS = [
|
||||
("wl-copy", ["wl-copy"], ["wl-paste"]),
|
||||
("xclip", ["xclip", "-selection", "clipboard"], ["xclip", "-selection", "clipboard", "-o"]),
|
||||
("pbcopy", ["pbcopy"], ["pbpaste"]),
|
||||
("xsel", ["xsel", "--clipboard", "--input"], ["xsel", "--clipboard", "--output"]),
|
||||
]
|
||||
|
||||
|
||||
def _backend(wanted: str | None) -> tuple[str, list[str], list[str]]:
|
||||
if wanted and wanted != "auto":
|
||||
for name, _copy, _paste in BACKENDS:
|
||||
if name == wanted:
|
||||
return name, _copy, _paste
|
||||
raise RuntimeError(f"unknown clipboard backend: {wanted}")
|
||||
for name, _copy, _paste in BACKENDS:
|
||||
if shutil.which(_copy[0]) and shutil.which(_paste[0]):
|
||||
return name, _copy, _paste
|
||||
raise RuntimeError("no clipboard backend found (install xclip or wl-copy/wl-paste)")
|
||||
|
||||
|
||||
def backend_available(wanted: str | None = None) -> bool:
|
||||
try:
|
||||
_backend(wanted)
|
||||
return True
|
||||
except RuntimeError:
|
||||
return False
|
||||
|
||||
|
||||
def copy(text: str, wanted: str | None = None) -> bool:
|
||||
try:
|
||||
name, cmd, _paste = _backend(wanted)
|
||||
except RuntimeError as exc:
|
||||
sys.stderr.write(f"shellbound: copy: {exc}\n")
|
||||
return False
|
||||
try:
|
||||
subprocess.run(
|
||||
cmd, input=text.encode(), check=True, stdout=subprocess.DEVNULL
|
||||
)
|
||||
return True
|
||||
except (OSError, subprocess.SubprocessError) as exc:
|
||||
sys.stderr.write(f"shellbound: copy via {name} failed: {exc}\n")
|
||||
return False
|
||||
|
||||
|
||||
def paste(wanted: str | None = None) -> str:
|
||||
try:
|
||||
name, _copy, cmd = _backend(wanted)
|
||||
except RuntimeError as exc:
|
||||
sys.stderr.write(f"shellbound: paste: {exc}\n")
|
||||
return ""
|
||||
try:
|
||||
return subprocess.run(cmd, check=True, capture_output=True).stdout.decode()
|
||||
except (OSError, subprocess.SubprocessError) as exc:
|
||||
sys.stderr.write(f"shellbound: paste via {name} failed: {exc}\n")
|
||||
return ""
|
||||
@@ -0,0 +1,93 @@
|
||||
"""Shellbound configuration.
|
||||
|
||||
Stored in ``~/.shellbound/config.json``. Environment variables override the
|
||||
file values:
|
||||
|
||||
- ``SHELLBOUND_GATEWAY`` gateway URL
|
||||
- ``SHELLBOUND_TOKEN`` API key
|
||||
- ``SHELLBOUND_VERIFY_TLS`` "true"/"1" to verify TLS certificates
|
||||
- ``SHELLBOUND_KEEP_WORKSPACES`` how many shellbound_* workspaces to keep open
|
||||
- ``SHELLBOUND_CLIPBOARD`` clipboard backend (auto/xclip/wl-copy/...)
|
||||
|
||||
``node_base_template`` was removed from the schema (never changes); any stale
|
||||
value left in an existing config is pruned on load.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
CONFIG_DIR = Path.home() / ".shellbound"
|
||||
CONFIG_FILE = CONFIG_DIR / "config.json"
|
||||
|
||||
DEFAULT_GATEWAY = "https://soda.lab.obelisk.cc:8443"
|
||||
DEFAULT_KEEP_WORKSPACES = 5
|
||||
|
||||
|
||||
class Config(dict):
|
||||
"""Dict-backed config with attribute access for convenience."""
|
||||
|
||||
def __getattr__(self, item: str):
|
||||
if item.startswith("_"):
|
||||
raise AttributeError(item)
|
||||
try:
|
||||
return self[item]
|
||||
except KeyError:
|
||||
raise AttributeError(item) from None
|
||||
|
||||
def __setattr__(self, item: str, value) -> None:
|
||||
self[item] = value
|
||||
|
||||
@classmethod
|
||||
def load(cls) -> "Config":
|
||||
data: dict = {}
|
||||
if CONFIG_FILE.exists():
|
||||
try:
|
||||
data = json.loads(CONFIG_FILE.read_text())
|
||||
except (json.JSONDecodeError, OSError):
|
||||
data = {}
|
||||
data.pop("node_base_template", None)
|
||||
env = {
|
||||
"SHELLBOUND_GATEWAY": "gateway",
|
||||
"SHELLBOUND_TOKEN": "token",
|
||||
"SHELLBOUND_VERIFY_TLS": "verify_tls",
|
||||
"SHELLBOUND_KEEP_WORKSPACES": "keep_workspaces",
|
||||
"SHELLBOUND_CLIPBOARD": "clipboard",
|
||||
}
|
||||
for env_key, cfg_key in env.items():
|
||||
val = os.environ.get(env_key)
|
||||
if val is None:
|
||||
continue
|
||||
if cfg_key == "verify_tls":
|
||||
data[cfg_key] = val.lower() in ("1", "true", "yes", "on")
|
||||
elif cfg_key == "keep_workspaces":
|
||||
try:
|
||||
data[cfg_key] = int(val)
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
data[cfg_key] = val
|
||||
data.setdefault("gateway", DEFAULT_GATEWAY)
|
||||
data.setdefault("token", "")
|
||||
data.setdefault("verify_tls", False)
|
||||
data.setdefault("clipboard", "auto")
|
||||
data.setdefault("keep_workspaces", DEFAULT_KEEP_WORKSPACES)
|
||||
return cls(data)
|
||||
|
||||
def save(self) -> None:
|
||||
CONFIG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
CONFIG_FILE.write_text(json.dumps(dict(self), indent=2) + "\n")
|
||||
|
||||
@property
|
||||
def has_token(self) -> bool:
|
||||
return bool(self.get("token"))
|
||||
|
||||
@property
|
||||
def keep_workspaces(self) -> int:
|
||||
try:
|
||||
value = int(self.get("keep_workspaces", DEFAULT_KEEP_WORKSPACES))
|
||||
except (TypeError, ValueError):
|
||||
return DEFAULT_KEEP_WORKSPACES
|
||||
return max(value, 0)
|
||||
@@ -0,0 +1,88 @@
|
||||
"""Mode-specific behaviors: expert two-stage orchestration, shell command
|
||||
extraction and execution, and small helpers.
|
||||
|
||||
Expert mode runs two stages against the ``shellbound_expert`` persona (whose
|
||||
system prompt already carries the expert-brief guidance, so none of that detail
|
||||
needs to be injected into the conversation). Only short stage markers plus the
|
||||
relevant payload are sent as turn messages.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
STAGE_1_TASK = (
|
||||
"Stage 1 - Produce an Expert Brief to answer the corresponding user prompt.\n"
|
||||
"Produce the Expert Brief only.\n\n"
|
||||
"USER PROMPT:\n{user_query}"
|
||||
)
|
||||
|
||||
STAGE_2_TASK = (
|
||||
"Stage 2 - Adopt the following Expert Persona to answer the user prompt.\n"
|
||||
"Follow the Expert Persona exactly, then answer the user prompt.\n\n"
|
||||
"EXPERT PERSONA (Expert Brief):\n{meta}\n\n"
|
||||
"USER PROMPT:\n{user_query}"
|
||||
)
|
||||
|
||||
|
||||
def expert_brief_task(user_query: str) -> str:
|
||||
return STAGE_1_TASK.format(user_query=user_query)
|
||||
|
||||
|
||||
def expert_retry_task(user_query: str, meta: str) -> str:
|
||||
return STAGE_2_TASK.format(meta=meta, user_query=user_query)
|
||||
|
||||
|
||||
def session_name(query: str, limit: int = 60) -> str:
|
||||
flat = " ".join(query.split())
|
||||
return f"sb: {flat[:limit]}"
|
||||
|
||||
|
||||
def minimal_uid(ws_id: str, others: list[str]) -> str:
|
||||
"""Shortest prefix of ``ws_id`` that is unique among ``others``."""
|
||||
for size in range(4, len(ws_id) + 1):
|
||||
prefix = ws_id[:size]
|
||||
if not any(other != ws_id and other.startswith(prefix) for other in others):
|
||||
return prefix
|
||||
return ws_id[:12]
|
||||
|
||||
|
||||
def extract_commands(raw: str) -> list[str]:
|
||||
"""Extract shell commands from a response.
|
||||
|
||||
Prefers the content of a fenced code block; otherwise uses non-empty,
|
||||
non-comment lines. Shared lines are stripped of a leading ``$``.
|
||||
"""
|
||||
text = (raw or "").strip()
|
||||
if not text:
|
||||
return []
|
||||
fenced: list[str] = []
|
||||
in_fence = False
|
||||
for line in text.splitlines():
|
||||
s = line.strip()
|
||||
if s.startswith("```") or s == "~~~":
|
||||
in_fence = not in_fence
|
||||
continue
|
||||
if in_fence:
|
||||
fenced.append(line)
|
||||
source = fenced if fenced else text.splitlines()
|
||||
commands: list[str] = []
|
||||
for line in source:
|
||||
s = line.strip()
|
||||
if not s or s.startswith("#"):
|
||||
continue
|
||||
if s.startswith("$"):
|
||||
s = s[1:].strip()
|
||||
commands.append(s)
|
||||
return commands
|
||||
|
||||
|
||||
def run_commands(commands: list[str]) -> int:
|
||||
"""Run each proposed command locally through ``bash -c`` (top-level)."""
|
||||
for command in commands:
|
||||
sys.stderr.write(f"$ {command}\n")
|
||||
rc = subprocess.call(["bash", "-c", command])
|
||||
if rc != 0:
|
||||
sys.stderr.write(f"shellbound: command exited with status {rc}\n")
|
||||
return 0
|
||||
@@ -0,0 +1,100 @@
|
||||
"""Persona (mode) resolution.
|
||||
|
||||
Modes are Turnstone personae whose slugs start with ``shellbound_``. A mode
|
||||
argument is resolved as a case-insensitive shortest-unique prefix of that
|
||||
slug. Zero matches means the argument is treated as part of the query; more
|
||||
than one match (e.g. ``e`` for ``explain``/``expert``) is an ambiguity error.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Optional
|
||||
|
||||
SHELLBOUND_PREFIX = "shellbound_"
|
||||
|
||||
# Canonical modes this CLI knows about, used to detect "persona not yet
|
||||
# provisioned on the server" so intent is never silently swallowed into the
|
||||
# query.
|
||||
KNOWN_MODES = {
|
||||
"shell": "shellbound_shell",
|
||||
"s": "shellbound_shell",
|
||||
"answer": "shellbound_answer",
|
||||
"a": "shellbound_answer",
|
||||
"explain": "shellbound_explain",
|
||||
"creative": "shellbound_creative",
|
||||
"c": "shellbound_creative",
|
||||
"poetic": "shellbound_poetic",
|
||||
"p": "shellbound_poetic",
|
||||
"expert": "shellbound_expert",
|
||||
}
|
||||
|
||||
|
||||
class AmbiguousPersona(ValueError):
|
||||
def __init__(self, fragment: str, candidates: list[str]):
|
||||
self.fragment = fragment
|
||||
self.candidates = candidates
|
||||
super().__init__(
|
||||
f"'{fragment}' is ambiguous ({len(candidates)} personae: "
|
||||
f"{', '.join(sorted(candidates))}); supply more characters"
|
||||
)
|
||||
|
||||
|
||||
def shellbound_personas(raw: list[dict]) -> list[dict]:
|
||||
return [p for p in raw if p.get("name", "").startswith(SHELLBOUND_PREFIX)]
|
||||
|
||||
|
||||
def interactions(raw: list[dict]) -> list[dict]:
|
||||
return [p for p in raw if "interactive" in (p.get("applies_to_kinds") or [])]
|
||||
|
||||
|
||||
def can_expected_mode(fragment: str) -> bool:
|
||||
return fragment.strip().lower() in KNOWN_MODES
|
||||
|
||||
|
||||
def resolve_mode(raw: list[dict], fragment: str) -> Optional[dict]:
|
||||
"""Resolve a mode fragment to a ``shellbound_*`` persona.
|
||||
|
||||
Returns ``None`` when the fragment matches nothing (it should then be
|
||||
treated as the start of the query). Raises :class:`AmbiguousPersona`
|
||||
when more than one persona matches.
|
||||
"""
|
||||
fragment = fragment.strip()
|
||||
if not fragment:
|
||||
return None
|
||||
low = fragment.lower()
|
||||
matches: list[dict] = []
|
||||
for p in shellbound_personas(raw):
|
||||
name = p["name"]
|
||||
if name.lower() == low or name.lower().startswith(SHELLBOUND_PREFIX + low):
|
||||
matches.append(p)
|
||||
if len(matches) > 1:
|
||||
raise AmbiguousPersona(fragment, [p["name"] for p in matches])
|
||||
return matches[0] if matches else None
|
||||
|
||||
|
||||
def resolve_name(raw: list[dict], name: str) -> dict:
|
||||
"""Resolve an explicit ``--persona`` value (exact name or unique prefix)."""
|
||||
name = name.strip()
|
||||
if not name:
|
||||
raise ValueError("empty persona name")
|
||||
low = name.lower()
|
||||
matches = [
|
||||
p for p in interactions(raw)
|
||||
if p["name"].lower() == low or p["name"].lower().startswith(low)
|
||||
]
|
||||
if not matches:
|
||||
available = ", ".join(sorted(p["name"] for p in raw))
|
||||
raise ValueError(f"unknown persona '{name}' (available: {available})")
|
||||
if len(matches) > 1:
|
||||
raise AmbiguousPersona(name, [p["name"] for p in matches])
|
||||
return matches[0]
|
||||
|
||||
|
||||
def default_persona_name(raw: list[dict]) -> Optional[str]:
|
||||
"""Preferred default persona (answer) with fallbacks."""
|
||||
sl = [p["name"] for p in shellbound_personas(raw)]
|
||||
for candidate in ("shellbound_answer", "shellbound_shell"):
|
||||
if candidate in sl:
|
||||
return candidate
|
||||
default = next((p["name"] for p in raw if p.get("is_default")), None)
|
||||
return default
|
||||
@@ -0,0 +1,184 @@
|
||||
"""Streaming terminal renderer.
|
||||
|
||||
In a TTY the assistant's reply renders live into a :class:`rich.live.Live`
|
||||
frame: already-completed Markdown blocks are laid out once (boundary-flushed)
|
||||
while the in-fight tail streams as bright-white plain text, and the thinking
|
||||
shows as dim italic text above. On a non-TTY or with ``--plain`` the raw
|
||||
Markdown source streams to stdout instead (reasoning stays silent there).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sys
|
||||
from typing import Optional
|
||||
|
||||
from rich.console import Console, Group, RenderableType
|
||||
from rich.live import Live
|
||||
from rich.markdown import Markdown
|
||||
from rich.rule import Rule
|
||||
from rich.text import Text
|
||||
|
||||
REASONING_STYLE = "bright_black italic"
|
||||
CONTENT_STYLE = "white"
|
||||
WARN_STYLE = "yellow"
|
||||
ERROR_STYLE = "red bold"
|
||||
FORCE_FLUSH_CHARS = 4000
|
||||
|
||||
|
||||
def is_tty() -> bool:
|
||||
return sys.stdout.isatty()
|
||||
|
||||
|
||||
class StreamRenderer:
|
||||
"""Accumulates a turn and renders it live (or plainly)."""
|
||||
|
||||
def __init__(self, *, plain: Optional[bool] = None):
|
||||
auto_plain = not is_tty()
|
||||
self.plain = plain if plain is not None else auto_plain
|
||||
self.reasoning: list[str] = []
|
||||
self.content_parts: list[str] = []
|
||||
self.flushed: list[RenderableType] = []
|
||||
self.tail = ""
|
||||
self._messages: list[RenderableType] = [] # sticky info/error lines
|
||||
self._live: Optional[Live] = None
|
||||
self._rendered = 0
|
||||
self._console = Console(
|
||||
color_system="auto",
|
||||
markup=False,
|
||||
highlight=None,
|
||||
tab_size=4,
|
||||
force_terminal=not self.plain and is_tty(),
|
||||
)
|
||||
|
||||
# -- lifecycle -------------------------------------------------------
|
||||
|
||||
def start(self) -> None:
|
||||
if not self.plain:
|
||||
self._live = Live(
|
||||
console=self._console,
|
||||
refresh_per_second=12,
|
||||
vertical_overflow="visible",
|
||||
)
|
||||
self._live.start(refresh=False)
|
||||
|
||||
def stop(self) -> None:
|
||||
if self._live is not None:
|
||||
self._live.stop()
|
||||
self._live = None
|
||||
|
||||
def __enter__(self) -> "StreamRenderer":
|
||||
self.start()
|
||||
return self
|
||||
|
||||
def __exit__(self, *exc) -> None:
|
||||
if not self.plain:
|
||||
self.finish()
|
||||
self.stop()
|
||||
|
||||
# -- content ---------------------------------------------------------
|
||||
|
||||
def add_reasoning(self, text: str) -> None:
|
||||
if not text:
|
||||
return
|
||||
self.reasoning.append(text)
|
||||
if not self.plain:
|
||||
self.refresh()
|
||||
|
||||
def add_content(self, text: str) -> None:
|
||||
if not text:
|
||||
return
|
||||
self.content_parts.append(text)
|
||||
if self.plain:
|
||||
sys.stdout.write(text)
|
||||
sys.stdout.flush()
|
||||
return
|
||||
self.tail += text
|
||||
self._maybe_flush()
|
||||
self.refresh()
|
||||
|
||||
def add_info(self, message: str) -> None:
|
||||
if not message:
|
||||
return
|
||||
if self.plain:
|
||||
sys.stderr.write(f"shellbound: {message}\n")
|
||||
sys.stderr.flush()
|
||||
else:
|
||||
self._messages.append(Text(str(message), style=WARN_STYLE))
|
||||
self.refresh()
|
||||
|
||||
def add_error(self, message: str) -> None:
|
||||
if self.plain:
|
||||
sys.stderr.write(f"shellbound: {message}\n")
|
||||
sys.stderr.flush()
|
||||
else:
|
||||
self._messages.append(Text(f"error: {message}", style=ERROR_STYLE))
|
||||
self.refresh()
|
||||
|
||||
# -- flushing --------------------------------------------------------
|
||||
|
||||
def _fence_count(self) -> int:
|
||||
return self.tail.count("```")
|
||||
|
||||
def _maybe_flush(self) -> None:
|
||||
if not self.tail:
|
||||
return
|
||||
fences = self._fence_count()
|
||||
if fences % 2 == 1:
|
||||
# inside an open fenced block: only force-flush raw at the cap
|
||||
if len(self.tail) > FORCE_FLUSH_CHARS:
|
||||
self._flush(markdown=False)
|
||||
return
|
||||
if self.tail.rstrip().endswith("```"):
|
||||
# fenced block just closed
|
||||
self._flush(markdown=True)
|
||||
return
|
||||
if self.tail.endswith("\n\n"):
|
||||
self._flush(markdown=True)
|
||||
elif len(self.tail) > FORCE_FLUSH_CHARS:
|
||||
self._flush(markdown=False)
|
||||
|
||||
def _flush(self, *, markdown: bool) -> None:
|
||||
block = self.tail
|
||||
self.tail = ""
|
||||
if not block.strip():
|
||||
return
|
||||
if markdown:
|
||||
self.flushed.append(Markdown(block))
|
||||
else:
|
||||
self.flushed.append(Text(block, style=CONTENT_STYLE))
|
||||
|
||||
# -- display ---------------------------------------------------------
|
||||
|
||||
def _renderable(self) -> RenderableType:
|
||||
items: list[RenderableType] = []
|
||||
if self.reasoning:
|
||||
items.append(Text("".join(self.reasoning), style=REASONING_STYLE))
|
||||
if self.content_parts or self.flushed or self.tail:
|
||||
items.append(Rule(style="bright_black"))
|
||||
items.extend(self.flushed)
|
||||
if self.tail:
|
||||
items.append(Text(self.tail, style=CONTENT_STYLE))
|
||||
items.extend(self._messages)
|
||||
if not items:
|
||||
return Text("")
|
||||
if len(items) == 1:
|
||||
return items[0]
|
||||
return Group(*items)
|
||||
|
||||
def refresh(self) -> None:
|
||||
if self._live is not None:
|
||||
self._live.update(self._renderable())
|
||||
|
||||
def finish(self) -> str:
|
||||
"""Flush any pending tail and return the raw streamed reply text."""
|
||||
if not self.plain and self.tail.strip():
|
||||
self._flush(markdown=True)
|
||||
if not self.plain:
|
||||
self.refresh()
|
||||
return "".join(self.content_parts)
|
||||
|
||||
# -- plain helpers ---------------------------------------------------
|
||||
|
||||
@property
|
||||
def reasoning_text(self) -> str:
|
||||
return "".join(self.reasoning)
|
||||
@@ -0,0 +1,112 @@
|
||||
"""Interactive setup for shellbound.
|
||||
|
||||
Runs outside asyncio (``input``/``getpass``) and rewrites
|
||||
``~/.shellbound/config.json``. Existing values are remembered: pressing Enter
|
||||
keeps the current value. The API key is never echoed; it is entered masked via
|
||||
:func:`getpass`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import getpass
|
||||
import urllib.parse
|
||||
|
||||
from .config import Config
|
||||
|
||||
CLIPBOARD_BACKENDS = ("auto", "xclip", "xsel", "wl-copy", "pbcopy")
|
||||
|
||||
|
||||
def _mask(value: str, show: int = 4) -> str:
|
||||
if not value:
|
||||
return "(not set)"
|
||||
if len(value) <= show:
|
||||
return "•" * len(value)
|
||||
return "•" * (len(value) - show) + value[-show:]
|
||||
|
||||
|
||||
def _edit(question: str, current: str, display: str | None = None) -> str:
|
||||
shown = _mask(current) if display == "mask" else current or "(not set)"
|
||||
prompt = f"{question} [{shown}] > "
|
||||
try:
|
||||
raw = input(prompt)
|
||||
except (EOFError, KeyboardInterrupt):
|
||||
print()
|
||||
return current
|
||||
cleaned = raw.strip()
|
||||
return current if not cleaned else cleaned
|
||||
|
||||
|
||||
def _edit_bool(question: str, current: bool) -> bool:
|
||||
label = "yes" if current else "no"
|
||||
while True:
|
||||
value = _edit(question, label)
|
||||
low = value.lower()
|
||||
if low in ("1", "true", "yes", "y", "on"):
|
||||
return True
|
||||
if low in ("0", "false", "no", "n", "off"):
|
||||
return False
|
||||
print(" answer y/n")
|
||||
|
||||
|
||||
def _edit_clipboard(current: str) -> str:
|
||||
while True:
|
||||
value = _edit(
|
||||
f"Clipboard backend ({'/'.join(CLIPBOARD_BACKENDS)})",
|
||||
current,
|
||||
)
|
||||
if value.lower() not in CLIPBOARD_BACKENDS:
|
||||
print(f" choose one of: {', '.join(CLIPBOARD_BACKENDS)}")
|
||||
continue
|
||||
return value.lower()
|
||||
|
||||
|
||||
def _edit_token(current: str) -> str:
|
||||
prompt = getpass.getpass(
|
||||
f"API key [paste a new key, or Enter to keep {_mask(current)}] > "
|
||||
)
|
||||
cleaned = prompt.strip()
|
||||
return current if not cleaned else cleaned
|
||||
|
||||
|
||||
def setup(cfg: Config) -> Config:
|
||||
"""Interactive config walkthrough; returns the updated config."""
|
||||
print("shellbound setup - press Enter to keep the current value.\n")
|
||||
|
||||
gateway = _edit("Gateway", cfg.get("gateway", ""))
|
||||
parsed = urllib.parse.urlparse(gateway)
|
||||
while parsed.scheme not in ("http", "https") or not parsed.netloc:
|
||||
print(f" not a valid gateway URL: {gateway!r}")
|
||||
gateway = _edit("Gateway", cfg.get("gateway", ""))
|
||||
parsed = urllib.parse.urlparse(gateway)
|
||||
|
||||
token = _edit_token(cfg.get("token", ""))
|
||||
keep = _edit("Keep open this many recent shellbound_* workspaces", str(cfg.keep_workspaces))
|
||||
while True:
|
||||
try:
|
||||
keep = max(int(keep), 0)
|
||||
break
|
||||
except ValueError:
|
||||
print(f" invalid number: {keep!r}")
|
||||
keep = _edit(
|
||||
"Keep open this many recent shellbound_* workspaces", str(cfg.keep_workspaces)
|
||||
)
|
||||
verify = _edit_bool(
|
||||
"Verify TLS certificates (self-signed lab cert: choose no)",
|
||||
cfg.get("verify_tls", False),
|
||||
)
|
||||
clipboard = _edit_clipboard(cfg.get("clipboard", "auto"))
|
||||
|
||||
cfg["gateway"] = gateway
|
||||
cfg["token"] = token
|
||||
cfg["keep_workspaces"] = keep
|
||||
cfg["verify_tls"] = bool(verify)
|
||||
cfg["clipboard"] = clipboard
|
||||
cfg.save()
|
||||
|
||||
print("\nshellbound setup complete.")
|
||||
print(f" gateway: {gateway}")
|
||||
print(f" API key: {_mask(token)}")
|
||||
print(f" keep_workspaces: {keep}")
|
||||
print(f" verify_tls: {'yes' if verify else 'no'}")
|
||||
print(f" clipboard: {clipboard}")
|
||||
return cfg
|
||||
@@ -0,0 +1,234 @@
|
||||
"""Unit tests for shellbound (no network required)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
|
||||
from shellbound import modes
|
||||
from shellbound.personas import (
|
||||
AmbiguousPersona,
|
||||
can_expected_mode,
|
||||
default_persona_name,
|
||||
resolve_mode,
|
||||
resolve_name,
|
||||
shellbound_personas,
|
||||
)
|
||||
|
||||
RAW = [
|
||||
{"name": "shellbound_shell", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "shellbound_answer", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "shellbound_creative", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "shellbound_expert", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "shellbound_explain", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "shellbound_poetic", "applies_to_kinds": ["interactive"]},
|
||||
{"name": "engineer", "applies_to_kinds": ["interactive"], "is_default": True},
|
||||
{"name": "orchestrator", "applies_to_kinds": ["coordinator"]},
|
||||
]
|
||||
|
||||
|
||||
def test_shellbound_filter():
|
||||
slugs = {p["name"] for p in shellbound_personas(RAW)}
|
||||
assert slugs == {
|
||||
"shellbound_shell",
|
||||
"shellbound_answer",
|
||||
"shellbound_creative",
|
||||
"shellbound_expert",
|
||||
"shellbound_explain",
|
||||
"shellbound_poetic",
|
||||
}
|
||||
|
||||
|
||||
def test_unique_abbreviations():
|
||||
assert resolve_mode(RAW, "s")["name"] == "shellbound_shell"
|
||||
assert resolve_mode(RAW, "a")["name"] == "shellbound_answer"
|
||||
assert resolve_mode(RAW, "c")["name"] == "shellbound_creative"
|
||||
assert resolve_mode(RAW, "p")["name"] == "shellbound_poetic"
|
||||
assert resolve_mode(RAW, "sh")["name"] == "shellbound_shell"
|
||||
assert resolve_mode(RAW, "expe")["name"] == "shellbound_expert"
|
||||
assert resolve_mode(RAW, "expl")["name"] == "shellbound_explain"
|
||||
assert resolve_mode(RAW, "shellbound_shell")["name"] == "shellbound_shell"
|
||||
assert resolve_mode(RAW, "SHELLBOUND_ANSWER")["name"] == "shellbound_answer"
|
||||
|
||||
|
||||
def test_ambiguous_abbreviation():
|
||||
try:
|
||||
resolve_mode(RAW, "e")
|
||||
except AmbiguousPersona as exc:
|
||||
assert set(exc.candidates) == {"shellbound_expert", "shellbound_explain"}
|
||||
else:
|
||||
raise AssertionError("expected AmbiguousPersona")
|
||||
|
||||
|
||||
def test_no_match_is_query():
|
||||
assert resolve_mode(RAW, "zzz") is None
|
||||
assert can_expected_mode("zzz") is False
|
||||
assert can_expected_mode("s") is True
|
||||
|
||||
|
||||
def test_resolve_name():
|
||||
assert resolve_name(RAW, "engineer")["name"] == "engineer"
|
||||
assert resolve_name(RAW, "shellbound_poetic")["name"] == "shellbound_poetic"
|
||||
try:
|
||||
resolve_name(RAW, "orchestrator") # coordinator kind, not interactive
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError("orchestrator should be rejected (non-interactive)")
|
||||
|
||||
|
||||
def test_default_persona():
|
||||
assert default_persona_name(RAW) == "shellbound_answer"
|
||||
|
||||
|
||||
def test_extract_commands_fenced():
|
||||
raw = '```\nls -la\n# hide me\necho hi\n```\n'
|
||||
assert modes.extract_commands(raw) == ["ls -la", "echo hi"]
|
||||
|
||||
|
||||
def test_extract_commands_plain():
|
||||
assert modes.extract_commands("whoami\necho ok") == ["whoami", "echo ok"]
|
||||
assert modes.extract_commands("# comment only") == []
|
||||
|
||||
|
||||
def test_extract_commands_dollar_prompt():
|
||||
assert modes.extract_commands("$ git status") == ["git status"]
|
||||
|
||||
|
||||
def test_minimal_uid():
|
||||
assert modes.minimal_uid("abc12345", ["abc99999", "abc12346"]) == "abc12345"
|
||||
assert modes.minimal_uid("aaa11111", []) == "aaa1"
|
||||
|
||||
|
||||
def test_session_name():
|
||||
assert modes.session_name(" why is the sky blue? ") == "sb: why is the sky blue?"
|
||||
|
||||
|
||||
def test_expert_prompts():
|
||||
task = modes.expert_brief_task("hello")
|
||||
assert task.startswith("Stage 1 - Produce an Expert Brief")
|
||||
assert "USER PROMPT:" in task and "hello" in task
|
||||
assert "prompt writing agent" not in task
|
||||
retry = modes.expert_retry_task("hello", "META")
|
||||
assert retry.startswith("Stage 2 - Adopt the following Expert Persona")
|
||||
assert "EXPERT PERSONA" in retry and "META" in retry
|
||||
assert "USER PROMPT:" in retry and "hello" in retry
|
||||
|
||||
|
||||
def _ws(ws_id, persona, state, updated):
|
||||
ns = type("W", (), {})()
|
||||
ns.ws_id = ws_id
|
||||
ns.persona = persona
|
||||
ns.state = state
|
||||
ns.updated = updated
|
||||
return ns
|
||||
|
||||
|
||||
def test_select_stale_keeps_newest_n():
|
||||
from shellbound.client import select_stale_workspaces
|
||||
|
||||
workstreams = [
|
||||
_ws("aaaa", "shellbound_answer", "idle", "2026-09-05T10:00:00"),
|
||||
_ws("bbbb", "shellbound_answer", "idle", "2026-09-05T11:00:00"),
|
||||
_ws("cccc", "engineer", "idle", "2026-09-05T12:00:00"), # not shellbound_
|
||||
_ws("dddd", "shellbound_shell", "idle", "2026-09-05T09:00:00"),
|
||||
]
|
||||
stale = select_stale_workspaces(workstreams, keep=2)
|
||||
# newest two shellbound_* kept (bbbb, aaaa); only the oldest is stale
|
||||
assert [w.ws_id for w in stale] == ["dddd"]
|
||||
|
||||
|
||||
def test_select_stale_excludes_closed_and_exclude_id():
|
||||
from shellbound.client import select_stale_workspaces
|
||||
|
||||
workstreams = [
|
||||
_ws("aaaa", "shellbound_answer", "closed", "2026-09-05T10:00:00"),
|
||||
_ws("bbbb", "shellbound_answer", "idle", "2026-09-05T11:00:00"),
|
||||
_ws("cccc", "shellbound_shell", "idle", "2026-09-05T09:00:00"),
|
||||
_ws("dddd", "shellbound_shell", "idle", "2026-09-05T08:00:00"),
|
||||
]
|
||||
# bbbb is the newest but excluded (it's the session just used); only
|
||||
# cccc remains inside keep=1, so the stale set is just dddd.
|
||||
stale = select_stale_workspaces(workstreams, keep=1, exclude_id="bbbb")
|
||||
assert [w.ws_id for w in stale] == ["dddd"]
|
||||
|
||||
|
||||
def test_select_stale_keep_zero():
|
||||
from shellbound.client import select_stale_workspaces
|
||||
|
||||
workstreams = [_ws("aaaa", "shellbound_answer", "idle", "2026-09-05T10:00:00")]
|
||||
assert [w.ws_id for w in select_stale_workspaces(workstreams, keep=0)] == ["aaaa"]
|
||||
|
||||
|
||||
def test_config_prunes_node_base_template(tmp_path, monkeypatch):
|
||||
from shellbound import config as c
|
||||
|
||||
cfg_file = tmp_path / "config.json"
|
||||
cfg_file.write_text(
|
||||
json.dumps({"gateway": "https://x:1", "token": "t", "node_base_template": "{gateway}/bad"})
|
||||
)
|
||||
monkeypatch.setattr(c, "CONFIG_FILE", cfg_file)
|
||||
cfg = c.Config.load()
|
||||
assert "node_base_template" not in cfg
|
||||
assert cfg["gateway"] == "https://x:1"
|
||||
assert cfg.keep_workspaces == 5
|
||||
|
||||
|
||||
def test_config_keep_workspaces_environ(tmp_path, monkeypatch):
|
||||
from shellbound import config as c
|
||||
|
||||
monkeypatch.setattr(c, "CONFIG_FILE", tmp_path / "missing.json")
|
||||
monkeypatch.setenv("SHELLBOUND_KEEP_WORKSPACES", "3")
|
||||
assert c.Config.load().keep_workspaces == 3
|
||||
|
||||
|
||||
def test_config_keep_workspaces_edge_values():
|
||||
from shellbound.config import Config
|
||||
|
||||
assert Config({"keep_workspaces": "oops"}).keep_workspaces == 5
|
||||
assert Config({}).keep_workspaces == 5
|
||||
assert Config({"keep_workspaces": 0}).keep_workspaces == 0
|
||||
|
||||
|
||||
def test_normalize_session_and_close():
|
||||
from shellbound.cli import _normalize_session
|
||||
|
||||
assert _normalize_session(["--session", "abc", "hi"]) == ["--session=abc", "hi"]
|
||||
assert _normalize_session(["--session"]) == ["--session=__list__"]
|
||||
assert _normalize_session(["--close", "all"]) == ["--close=all"]
|
||||
assert _normalize_session(["--close", "--debug"]) == ["--close=__all__", "--debug"]
|
||||
|
||||
|
||||
def test_renderer_flush_paragraphs():
|
||||
from shellbound.render import StreamRenderer
|
||||
|
||||
renderer = StreamRenderer(plain=False)
|
||||
renderer.add_content("First paragraph.\n")
|
||||
assert not renderer.flushed # in-flight tail, not yet a block
|
||||
assert renderer.tail.endswith("\n")
|
||||
renderer.add_content("Second.\n\n")
|
||||
assert len(renderer.flushed) == 1 and not renderer.tail
|
||||
|
||||
|
||||
def test_renderer_flush_code_fence():
|
||||
from shellbound.render import StreamRenderer
|
||||
|
||||
renderer = StreamRenderer(plain=False)
|
||||
# blank lines inside an open fence must not flush; the closed fence
|
||||
# becomes a single Markdown block on the trailing blank line.
|
||||
renderer.add_content("```bash\necho a\n\n echo b \n```\n\n")
|
||||
assert len(renderer.flushed) == 1 and not renderer.tail
|
||||
# an unclosed fence cap is force-flushed as raw text
|
||||
renderer.add_content("```python\ndef f():\n")
|
||||
assert len(renderer.flushed) == 1 # still inside the fence, nothing flushed
|
||||
long = "x" * 5000 + "\n"
|
||||
renderer.add_content(long)
|
||||
assert len(renderer.flushed) == 2 and not renderer.tail
|
||||
|
||||
|
||||
def test_renderer_plain_content_streams():
|
||||
from shellbound.render import StreamRenderer
|
||||
|
||||
renderer = StreamRenderer(plain=True)
|
||||
renderer.add_content("just text")
|
||||
assert renderer.content_parts == ["just text"]
|
||||
assert renderer.finish() == "just text"
|
||||
Reference in New Issue
Block a user