Skip to content
Open
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: 1 addition & 1 deletion Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ WORKDIR /app
COPY pyproject.toml uv.lock README.md LICENSE entrypoint.sh ./
# Copy the full source and install
COPY src ./src
RUN uv sync --no-dev --no-editable && uv pip install "any-llm-sdk[gemini,xai]" && \
RUN uv sync --locked --no-dev --no-editable && \
chmod +x /app/entrypoint.sh

WORKDIR /workspace
Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ See the [Build docs](https://bub.build/docs/build/) for hook guides, packaging,
| `bub gateway` | Channel listener (Telegram, etc.) |
| `bub install` | Install or sync Bub plugin deps |
| `bub update` | Upgrade Bub plugin deps |
| `bub login openai` | OpenAI Codex OAuth |
| `bub login codex` | ChatGPT plan login |

Lines starting with `,` enter internal command mode (`,help`, `,skill name=my-skill`, `,fs.read path=README.md`).

Expand All @@ -131,7 +131,7 @@ Lines starting with `,` enter internal command mode (`,help`, `,skill name=my-sk
| Variable | Default | Description |
| --------------------------- | ---------------------------- | ---------------------------------------------------- |
| `BUB_MODEL` | `openrouter:openrouter/free` | Model identifier |
| `BUB_API_KEY` | — | Provider key (optional with `bub login openai`) |
| `BUB_API_KEY` | — | Provider key (optional with `bub login codex`) |
| `BUB_API_BASE` | — | Custom provider endpoint |
| `BUB_CLIENT_ARGS` | — | JSON object forwarded to the underlying model client |
| `BUB_COMPLETION_ARGS` | — | JSON object forwarded to each completion call |
Expand Down
3 changes: 1 addition & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,7 @@ dependencies = [
"pluggy>=1.6.0",
"inquirer-textual>=0.8.0",
"typer>=0.24.1",
"authlib>=1.7.2",
"any-llm-sdk[anthropic]>=1.22.1",
"republic>=0.6.0a1",
"rich>=14.3.4",
"prompt-toolkit>=3.0.52",
"python-telegram-bot>=22.7",
Expand Down
74 changes: 38 additions & 36 deletions src/bub/builtin/agent.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
"""Runtime engine to process prompts with any-llm-sdk."""
"""Runtime engine to process prompts with Republic."""

from __future__ import annotations

Expand All @@ -15,6 +15,7 @@
from pathlib import Path
from typing import Any

import republic
from loguru import logger

from bub.builtin.commands import strip_command_prefix, validate_command_prefix
Expand All @@ -25,6 +26,7 @@
from bub.builtin.settings import load_settings
from bub.envelope import field_of
from bub.framework import BubFramework
from bub.prompt import LegacyPrompt, UserContent, prompt_text, to_content
from bub.skills import discover_skills, render_skills_prompt
from bub.store import AsyncTapeStore, AsyncTapeStoreAdapter, InMemoryTapeStore, TapeStore, is_async_tape_store
from bub.streaming import AsyncStreamEvents, StreamEvent, StreamState
Expand All @@ -39,7 +41,7 @@


class Agent:
"""Agent that processes prompts using hooks, tools, tape, and any-llm-sdk."""
"""Agent that processes prompts using hooks, tools, tape, and Republic."""

def __init__(
self,
Expand Down Expand Up @@ -115,7 +117,7 @@ async def run_stream(
self,
*,
session_id: str,
prompt: str | list[dict],
prompt: list[UserContent] | LegacyPrompt,
state: TurnState | None = None,
model: str | None = None,
allowed_skills: Collection[str] | None = None,
Expand All @@ -127,8 +129,9 @@ async def run_stream(
Args:
session_id: Session identity within the workspace. A ``temp/`` prefix
prevents the turn's fork from merging back into its parent tape.
prompt: Text or multimodal content parts. Text beginning with this
agent's command prefix after stripping whitespace invokes a command.
prompt: User content (text, images, audio and video); legacy text and content blocks are converted.
A text-only prompt beginning with this agent's command prefix after
stripping whitespace invokes a command.
state: Mutable turn state. None loads state through framework hooks
using this agent's store; supplied state skips that loading.
The current agent is always bound into the state.
Expand All @@ -154,11 +157,12 @@ async def run_stream(
"gen_ai.conversation.id": session_id,
},
)
span.messages("gen_ai.input.messages", [{"role": "user", "content": prompt}])
stack = AsyncExitStack()
try:
with span.activate():
if not prompt:
content = to_content(prompt)
span.messages("gen_ai.input.messages", [republic.user(*content)])
if not any(content):
events = self._events_from_iterable([
StreamEvent("text", {"delta": "error: empty prompt"}),
StreamEvent("final", {"text": "error: empty prompt", "ok": False}),
Expand All @@ -180,7 +184,11 @@ async def run_stream(
tape.fork_tape(merge_back=not session_id.startswith("temp/"))
)
await tape.ensure_bootstrap_anchor()
command = strip_command_prefix(prompt, self.command_prefix) if isinstance(prompt, str) else None
command = (
strip_command_prefix(content[0], self.command_prefix)
if len(content) == 1 and isinstance(content[0], str)
else None
)
if command is not None:
result = await self._run_command(tape=tape, line=command)
events = self._events_from_iterable([
Expand All @@ -190,7 +198,7 @@ async def run_stream(
else:
events = await self._agent_loop(
tape=tape,
prompt=prompt,
prompt=content,
model=model,
allowed_skills=allowed_skills,
allowed_tools=allowed_tools,
Expand Down Expand Up @@ -283,18 +291,18 @@ async def _agent_loop(
self,
*,
tape: Tape,
prompt: str | list[dict],
prompt: list[UserContent],
model: str | None = None,
allowed_skills: Collection[str] | None = None,
allowed_tools: Collection[str] | None = None,
) -> AsyncStreamEvents:
next_prompt: str | list[dict] = prompt
next_prompt = prompt
display_model = model or self.settings.model
await tape.append_event(
"loop.start",
{
"model": display_model,
"prompt": prompt,
"prompt": _prompt_payload(prompt),
"allowed_skills": list(allowed_skills) if allowed_skills else None,
"allowed_tools": list(allowed_tools) if allowed_tools else None,
},
Expand All @@ -313,20 +321,23 @@ async def _agent_loop(
async def _stream_events_with_auto_handoff(
self,
tape: Tape,
prompt: str | list[dict],
prompt: list[UserContent],
state: StreamState,
model: str | None = None,
allowed_skills: Collection[str] | None = None,
allowed_tools: Collection[str] | None = None,
) -> AsyncGenerator[StreamEvent, None]:
auto_handoff_remaining = MAX_AUTO_HANDOFF_RETRIES
display_model = model or self.settings.model
next_prompt: str | list[dict] | None = prompt
next_prompt: list[UserContent] | None = prompt
for step in range(1, self.settings.max_steps + 1):
start = time.monotonic()
should_continue = False
logger.info("loop.step step={} tape={} model={}", step, tape.name, display_model)
await tape.append_event("loop.step.start", {"step": step, "prompt": next_prompt})
await tape.append_event(
"loop.step.start",
{"step": step, "prompt": None if next_prompt is None else _prompt_payload(next_prompt)},
)
try:
output = await self._run_once(
tape=tape,
Expand Down Expand Up @@ -407,7 +418,8 @@ async def _stream_events_with_auto_handoff(
)
return

next_prompt = await self.framework.continue_prompt(prompt=next_prompt, tape=tape, state=state)
continuation = await self.framework.continue_prompt(prompt=next_prompt, tape=tape, state=state)
next_prompt = None if continuation is None else [continuation]
await tape.append_event(
"loop.step",
{
Expand All @@ -433,24 +445,17 @@ async def _run_once(
self,
*,
tape: Tape,
prompt: str | list[dict] | None,
prompt: list[UserContent] | None,
model: str | None = None,
allowed_tools: Collection[str] | None = None,
allowed_skills: Collection[str] | None = None,
) -> AsyncStreamEvents:
if isinstance(prompt, str):
prompt_text = prompt
elif prompt is None:
prompt_text = ""
else:
prompt_text = _extract_text_from_parts(prompt)
if allowed_skills is not None:
allowed_skills = {name.casefold() for name in allowed_skills}
tape.context.state["allowed_skills"] = list(allowed_skills)
return await self._run_once_stream(
tape=tape,
prompt=prompt,
prompt_text=prompt_text,
model=model,
allowed_skills=allowed_skills,
tools=self._allowed_tools(tape, allowed_tools),
Expand Down Expand Up @@ -492,8 +497,7 @@ async def _run_once_stream(
self,
*,
tape: Tape,
prompt: str | list[dict] | None,
prompt_text: str,
prompt: list[UserContent] | None,
model: str | None,
allowed_skills: set[str] | None,
tools: list[Tool],
Expand All @@ -502,7 +506,7 @@ async def _run_once_stream(
tools, code_mode_prompt = await self._prepare_code_mode(tools, tape)
tools, deferred_prompt = await self._prepare_deferred_tools(tools, tape)
system_prompt = self._system_prompt(
prompt_text,
prompt or [],
state=tape.context.state,
allowed_skills=allowed_skills,
tools_prompt="\n\n".join(block for block in (code_mode_prompt, deferred_prompt) if block),
Expand All @@ -512,13 +516,11 @@ async def _run_once_stream(
model_tools_for_call = model_tools(tools)
if (span := current_span()) and span.recording:
span.set(**{
"gen_ai.tool.definitions": [
tool.to_schema()["function"] | {"type": "function"} for tool in model_tools_for_call
]
"gen_ai.tool.definitions": [tool.to_schema() | {"type": "function"} for tool in model_tools_for_call]
})
steering_inbox = self.framework.get_steering_inbox()
steering_envelopes = await steering_inbox.drain_messages(tape.context.state) if steering_inbox else []
steering_messages = list(
steering_messages: list[list[UserContent]] = list(
await asyncio.gather(*[
self.framework.build_prompt(
message, session_id=field_of(message, "session_id"), state=tape.context.state
Expand Down Expand Up @@ -566,7 +568,7 @@ async def _prepare_code_mode(self, tools: list[Tool], tape: Tape) -> tuple[list[

def _system_prompt(
self,
prompt: str,
prompt: list[UserContent],
state: TurnState,
allowed_skills: set[str] | None = None,
tools_prompt: str = "",
Expand All @@ -577,7 +579,7 @@ def _system_prompt(
if tools_prompt:
blocks.append(tools_prompt)
workspace = workspace_from_state(state)
if skills_prompt := self._load_skills_prompt(prompt, workspace, allowed_skills):
if skills_prompt := self._load_skills_prompt(prompt_text(prompt), workspace, allowed_skills):
blocks.append(skills_prompt)
return "\n\n".join(blocks)

Expand Down Expand Up @@ -616,6 +618,6 @@ def _parse_args(args_tokens: list[str]) -> Args:
return Args(positional=positional, kwargs=kwargs)


def _extract_text_from_parts(parts: list[dict]) -> str:
"""Extract text content from multimodal content parts."""
return "\n".join(p.get("text", "") for p in parts if p.get("type") == "text")
def _prompt_payload(content: list[UserContent]) -> object:
"""Serialize prompt content for tape events."""
return republic.user(*content).to_dict().get("content")
Loading
Loading