Compare commits

..
3 Commits
Author SHA1 Message Date
github-actions[bot] a97c64c67b chore: update clone traffic data [skip ci] 2026-08-14 07:19:51 +00:00
Jon Saad-Falcon 64333651d1 Delete tools/pearl-reference-oracle directory 2026-08-13 19:27:37 -07:00
Elliot Slusky 9bee016c82 fix(server): preserve grounded agent stream content (#736)
* fix(server): preserve grounded agent stream content

* test(server): cover active grounded stream path

* fix(server): retain agent stream bridge
2026-08-13 18:15:36 -07:00
7 changed files with 118 additions and 311 deletions
+1 -1
View File
@@ -1,7 +1,7 @@
{
"schemaVersion": 1,
"label": "Git Clones",
"message": "189,482",
"message": "190,252",
"color": "green",
"namedLogo": "git"
}
+4 -3
View File
@@ -1,6 +1,6 @@
{
"total_clones": 189482,
"last_updated": "2026-08-13T07:22:13Z",
"total_clones": 190252,
"last_updated": "2026-08-14T07:19:51Z",
"daily": {
"2026-03-27": 2189,
"2026-03-28": 1874,
@@ -140,6 +140,7 @@
"2026-08-09": 1076,
"2026-08-10": 1060,
"2026-08-11": 2182,
"2026-08-12": 641
"2026-08-12": 641,
"2026-08-13": 770
}
}
+8 -55
View File
@@ -246,62 +246,15 @@ class AgentStreamBridge:
{"results": tool_results_data},
)
# Stream content using real LLM token streaming via
# engine.stream_full() when the engine is available.
# ``agent.run()`` already produced the authoritative, grounded
# response. Do not call the engine again here: a second inference
# would not have the agent's system prompt, tool transcript, or
# other internal context and could therefore contradict the
# result reported by the agent events. Replay the final content
# in chunks so the OpenAI-compatible streaming response stays
# consistent with the completed agent run.
content = agent_result.content or ""
engine = getattr(self._agent, "_engine", None)
used_real_streaming = False
if engine is not None and hasattr(engine, "stream_full") and content:
# Re-stream using the engine for real token delivery.
# Build the same messages the agent used for its final turn.
try:
from openjarvis.core.types import Message as MsgType
from openjarvis.core.types import Role as RoleType
replay_messages = []
for m in self._request.messages:
role = (
RoleType(m.role)
if m.role in {r.value for r in RoleType}
else RoleType.USER
)
replay_messages.append(
MsgType(
role=role,
content=m.content or "",
name=m.name,
tool_call_id=m.tool_call_id,
)
)
async for sc in engine.stream_full(
replay_messages,
model=self._model,
):
if sc.content:
chunk = ChatCompletionChunk(
id=self._chunk_id,
model=self._model,
choices=[
StreamChoice(
delta=DeltaMessage(content=sc.content),
)
],
)
yield f"data: {chunk.model_dump_json()}\n\n"
used_real_streaming = True
except Exception as stream_exc:
import logging as _logging
_logger = _logging.getLogger("openjarvis.server")
_logger.warning(
"Real streaming failed, falling back to word replay: %s",
stream_exc,
)
# Fallback: word-by-word replay if real streaming was not used
if not used_real_streaming and content:
if content:
words = content.split(" ")
for i, word in enumerate(words):
token = word if i == 0 else " " + word
+41
View File
@@ -847,6 +847,47 @@ class TestIdentityPromptInjection:
assert len(system_msgs) == 1
assert system_msgs[0].content == "Be terse."
def test_stream_uses_grounded_agent_result_without_replay(self):
"""Regression for #734: web streaming emits the agent's final answer."""
from openjarvis.core.events import EventBus
captured: list = []
engine = _make_capturing_engine(captured)
agent = _make_agent(content="My name is Jarvis Prime.")
agent._tools = [object()]
agent._engine = engine
client = TestClient(
create_app(
engine,
"test-model",
agent=agent,
bus=EventBus(),
config=_identity_config(),
)
)
resp = client.post(
"/v1/chat/completions",
json={
"model": "test-model",
"messages": [{"role": "user", "content": "who are you?"}],
"stream": True,
},
)
assert resp.status_code == 200
streamed_content = ""
for line in resp.text.splitlines():
if not line.startswith("data: {"):
continue
payload = json.loads(line.removeprefix("data: "))
choices = payload.get("choices", [])
if choices and choices[0]["delta"].get("content"):
streamed_content += choices[0]["delta"]["content"]
assert streamed_content == "My name is Jarvis Prime."
assert captured == []
agent.run.assert_called_once()
def test_direct_injects_identity_when_absent(self):
captured: list = []
engine = _make_capturing_engine(captured)
+64 -2
View File
@@ -1,6 +1,68 @@
import json
"""Regression tests for streaming completed agent responses."""
from openjarvis.server.stream_bridge import AgentStreamBridge
from __future__ import annotations
import asyncio
import json
from unittest.mock import MagicMock
import pytest
pytest.importorskip("fastapi")
from openjarvis.agents._stubs import AgentResult # noqa: E402
from openjarvis.core.events import EventBus # noqa: E402
from openjarvis.core.types import ToolResult # noqa: E402
from openjarvis.server.models import ChatCompletionRequest # noqa: E402
from openjarvis.server.stream_bridge import AgentStreamBridge # noqa: E402
def _streamed_content(events: list[str]) -> str:
"""Join assistant content from OpenAI-compatible data chunks."""
content = []
for event in events:
if not event.startswith("data: {"):
continue
payload = json.loads(event.removeprefix("data: ").strip())
choices = payload.get("choices")
if choices and choices[0]["delta"].get("content"):
content.append(choices[0]["delta"]["content"])
return "".join(content)
def test_stream_replays_grounded_agent_result_without_second_inference():
grounded_content = "My name is Jarvis. The tool reports 72 degrees."
agent = MagicMock()
agent._model = "configured-model"
agent.run.return_value = AgentResult(
content=grounded_content,
tool_results=[
ToolResult(tool_name="weather", content="72 degrees", success=True)
],
metadata={"prompt_tokens": 10, "completion_tokens": 12, "total_tokens": 22},
)
async def ungrounded_replay(*args, **kwargs):
raise AssertionError("stream_full must not run after agent.run")
yield # pragma: no cover
agent._engine.stream_full = ungrounded_replay
request = ChatCompletionRequest(
model="requested-model",
messages=[{"role": "user", "content": "Who are you, and what's outside?"}],
stream=True,
)
bridge = AgentStreamBridge(agent, EventBus(), request.model, request)
async def collect_events() -> list[str]:
return [event async for event in bridge.stream()]
events = asyncio.run(collect_events())
assert _streamed_content(events) == grounded_content
assert any(event.startswith("event: tool_results\n") for event in events)
agent.run.assert_called_once()
assert agent._model == "configured-model"
def test_tool_call_start_serializes_arguments_for_sse_without_mutating_event():
-111
View File
@@ -1,111 +0,0 @@
# Pearl reference oracle (OpenJarvis Phase 0 deliverable)
Phase 0-B of [Spec B](../../docs/design/2026-05-05-apple-silicon-pearl-mining-design.md)
called for "build a Python reference oracle for NoisyGEMM, validate against the
Pearl CUDA reference."
**Phase 0 found the oracle already exists upstream**, in two complementary forms:
| Layer | Upstream location | What it covers |
|---|---|---|
| Pure-Rust mining algorithm exposed to Python | `pearl/py-pearl-mining` | The complete `mine()` + `verify_plain_proof()` cycle. CPU-only. Hardware-portable. |
| PyTorch reference of production NoisyGEMM | `pearl/miner/miner-base/src/miner_base/noisy_gemm.py` | The same NoisyGEMM that vllm-miner accelerates with H100 CUDA. Bit-exact denoising verified by upstream test (`tests/test_noisy_gemm.py:92`). |
So this directory contains:
1. `smoke_test.py` — a runnable script that **actually mines a block on this machine** using the upstream Rust path, demonstrating the v1 architecture works on Apple Silicon (or any platform where `py-pearl-mining` builds).
2. This README documenting where the reference math lives.
## What this is *not*
This is **not a reimplementation** of NoisyGEMM. The original Spec B planned for that;
Phase 0 made it unnecessary. If you're tempted to write `noisy_gemm.py` here, stop —
read `pearl/miner/miner-base/src/miner_base/noisy_gemm.py` instead.
## Setup
You need:
- macOS arm64 (M1/M2/M3/M4) **or** Linux x86_64 / aarch64
- Python 3.12 (`uv venv --python 3.12 .venv` is the easiest)
- Rust 1.78+ (any recent toolchain — verified with 1.94 on macOS arm64)
- The Pearl source tree somewhere on disk
Build the wheel and install it (one-time, ~60 s on a fast Mac, ~5 min on first build):
```bash
# from the Pearl repo root
cd py-pearl-mining
uv pip install maturin
maturin build --release --interpreter "$(which python)"
# install the resulting wheel
uv pip install target/wheels/py_pearl_mining-*.whl
```
Or if Pearl publishes to PyPI in the future:
```bash
uv pip install py-pearl-mining
```
## Run the smoke test
```bash
python smoke_test.py
```
Actual output on Apple Silicon M2 Max (numbers will vary by hardware and run):
```
host: macOS-26.4.1-arm64-arm-64bit (arm64)
python: 3.12.1
[ok] pearl_mining loaded from <site-packages>/pearl_mining/__init__.py
[ok] PUBLICDATA_SIZE=164 MERKLE_LEAF_SIZE=1024
[ok] mine(m=256, n=128, k=1024, rank=32) returned a proof in 0.119 s
proof.m=256 proof.n=128 proof.k=1024 noise_rank=32
a.row_indices=[177, 185, 241, 249] bt.row_indices=[80, 81, 88, 89, 112, 113, 120, 121]
[ok] verify_plain_proof: ok=True ('Mining solution verified successfully', 0.2 ms)
[ok] all checks passed — Pearl mining works on this host
```
The `a.row_indices` and `bt.row_indices` values above are not constants — they're
`(offset + ROWS_PATTERN)` and `(offset + COLS_PATTERN)` for whichever offset the
miner happened to find a jackpot at. The smoke test verifies the *deltas* match
the configured `PeriodicPattern`, not the absolute values.
If it succeeds, this host can mine Pearl using the OpenJarvis `cpu-pearl` provider
(see Spec B §13). If it fails, the `[fail]` line tells you which step broke.
## What this proves (and what it doesn't)
**Proves:**
- The Pearl mining algorithm executes correctly on this host's CPU.
- Generated proofs verify under `verify_plain_proof`. (This is the same check
validators run on the inputs to the ZK proof.)
- The whole stack — `pearl-blake3`, `zk-pow`, `py-pearl-mining` — builds and
loads as a native CPython extension.
**Does NOT prove:**
- Network-difficulty hashrate. The smoke test uses
`nbits=0x1D2FFFFF` (test difficulty), much easier than mainnet. Real mining
expected hashrate on Apple Silicon CPU is several orders of magnitude lower
per share — see Spec B §1.5.6.
- ZK proof generation throughput. The smoke test calls `verify_plain_proof`,
not `generate_proof`. Plonky2 STARK proving takes seconds-to-minutes of CPU
per block (Spec B Open Q10).
- That this host can keep up with the network's block production rate.
## When to update this
- When Pearl bumps `py-pearl-mining` API: re-run the smoke test against the
new ref pinned in `OpenJarvis/src/openjarvis/mining/_constants.py`.
- When Pearl publishes a Mac wheel to PyPI: simplify the install instructions
above, drop the local `maturin build` step.
- When Spec B v2 adds the PyTorch-MPS reference path: extend `smoke_test.py`
with an MPS path comparison. The `miner-base` reference is already in
PyTorch, so the v2 smoke test would be a different test invoking
`miner_base.NoisyGemm` and comparing CPU vs MPS outputs for parity.
-139
View File
@@ -1,139 +0,0 @@
"""Pearl mining smoke test — runs an end-to-end mine + verify cycle.
Verifies that this host can run Pearl's pure-Rust mining algorithm via the
`pearl_mining` Python package. Used as Phase 0-B of the OpenJarvis Apple Silicon
mining spec ([Spec B]).
Exit codes:
0 all checks passed
1 pearl_mining import failed
2 mine() failed
3 verify_plain_proof rejected the proof
4 timing or sanity check failed
[Spec B]: ../../docs/design/2026-05-05-apple-silicon-pearl-mining-design.md
"""
from __future__ import annotations
import platform
import sys
import time
# Test fixture values — match upstream Pearl's tests/test_python_api.py so we are
# testing the same code path that Pearl's own CI exercises. Do not change
# without re-syncing with upstream.
DEFAULT_NBITS = 0x1D2FFFFF
DEFAULT_M = 256
DEFAULT_N = 128
DEFAULT_K = 1024
DEFAULT_RANK = 32
ROWS_PATTERN = [0, 8, 64, 72]
COLS_PATTERN = [0, 1, 8, 9, 32, 33, 40, 41]
def _ok(msg: str) -> None:
print(f"[ok] {msg}")
def _fail(msg: str, code: int) -> None:
print(f"[fail] {msg}")
sys.exit(code)
def main() -> None:
print(f"host: {platform.platform()} ({platform.machine()})")
print(f"python: {sys.version.split()[0]}")
try:
import pearl_mining
except ImportError as e:
_fail(f"could not import pearl_mining — install with `uv pip install py-pearl-mining` or build from source: {e}", 1)
_ok(f"pearl_mining loaded from {pearl_mining.__file__}")
_ok(
f"PUBLICDATA_SIZE={pearl_mining.PUBLICDATA_SIZE} "
f"MERKLE_LEAF_SIZE={pearl_mining.MERKLE_LEAF_SIZE}"
)
block_header = pearl_mining.IncompleteBlockHeader(
version=0,
prev_block=b"\x00" * 32,
merkle_root=b"0123456789abcdef" * 2,
timestamp=0x66666666,
nbits=DEFAULT_NBITS,
)
mining_config = pearl_mining.MiningConfiguration(
common_dim=DEFAULT_K,
rank=DEFAULT_RANK,
mma_type=pearl_mining.MMAType.Int7xInt7ToInt32,
rows_pattern=pearl_mining.PeriodicPattern.from_list(ROWS_PATTERN),
cols_pattern=pearl_mining.PeriodicPattern.from_list(COLS_PATTERN),
reserved=pearl_mining.MiningConfiguration.RESERVED,
)
t0 = time.perf_counter()
try:
plain_proof = pearl_mining.mine(
DEFAULT_M,
DEFAULT_N,
DEFAULT_K,
block_header,
mining_config,
signal_range=None,
wrong_jackpot_hash=False,
)
except Exception as e:
_fail(f"mine() raised: {e!r}", 2)
t_mine = time.perf_counter() - t0
_ok(
f"mine(m={DEFAULT_M}, n={DEFAULT_N}, k={DEFAULT_K}, rank={DEFAULT_RANK}) "
f"returned a proof in {t_mine:.3f} s"
)
print(
f" proof.m={plain_proof.m} proof.n={plain_proof.n} proof.k={plain_proof.k} "
f"noise_rank={plain_proof.noise_rank}"
)
print(
f" a.row_indices={plain_proof.a.row_indices} "
f"bt.row_indices={plain_proof.bt.row_indices}"
)
t0 = time.perf_counter()
ok, msg = pearl_mining.verify_plain_proof(block_header, plain_proof)
t_verify_ms = (time.perf_counter() - t0) * 1000
if not ok:
_fail(f"verify_plain_proof rejected our proof: {msg}", 3)
_ok(f"verify_plain_proof: ok=True ({msg!r}, {t_verify_ms:.1f} ms)")
if plain_proof.m != DEFAULT_M or plain_proof.n != DEFAULT_N or plain_proof.k != DEFAULT_K:
_fail("plain_proof dimensions do not match request", 4)
if plain_proof.noise_rank != DEFAULT_RANK:
_fail("plain_proof noise_rank does not match request", 4)
# Row indices are (offset + base_index) for some valid offset within the
# matrix dimension — see threads_partition() in zk-pow/src/ffi/mine.rs.
# We can't assert an absolute value (different offsets are valid every run),
# but we can assert the deltas match the pattern shape.
a_idxs = list(plain_proof.a.row_indices)
bt_idxs = list(plain_proof.bt.row_indices)
a_deltas = [v - a_idxs[0] for v in a_idxs]
bt_deltas = [v - bt_idxs[0] for v in bt_idxs]
if a_deltas != ROWS_PATTERN:
_fail(f"a.row_indices deltas ({a_deltas}) != ROWS_PATTERN ({ROWS_PATTERN})", 4)
if bt_deltas != COLS_PATTERN:
_fail(f"bt.row_indices deltas ({bt_deltas}) != COLS_PATTERN ({COLS_PATTERN})", 4)
print()
print("[ok] all checks passed — Pearl mining works on this host")
print()
print("Note: this used test difficulty (nbits=0x1D2FFFFF), not mainnet.")
print("Real-network shares per second will be many orders of magnitude lower.")
print("See docs/design/2026-05-05-apple-silicon-pearl-mining-design.md §1.5.6")
if __name__ == "__main__":
main()