Step 387: add orchestrator event stream APIs

This commit is contained in:
Bill
2026-02-16 13:01:15 -07:00
parent 67e637b9e8
commit 73e64295f2
9 changed files with 438 additions and 1 deletions

View File

@@ -2362,4 +2362,13 @@ target_link_libraries(step386_test PRIVATE
tree_sitter_javascript tree_sitter_typescript
tree_sitter_java tree_sitter_rust tree_sitter_go)
add_executable(step387_test tests/step387_test.cpp)
target_include_directories(step387_test PRIVATE src)
target_link_libraries(step387_test PRIVATE
nlohmann_json::nlohmann_json
unofficial::tree-sitter::tree-sitter
tree_sitter_python tree_sitter_cpp tree_sitter_elisp
tree_sitter_javascript tree_sitter_typescript
tree_sitter_java tree_sitter_rust tree_sitter_go)
# Step 12: Dear ImGui shell scaffolding created (main.cpp exists but not built due to dependencies)

View File

@@ -72,7 +72,9 @@ struct AgentPermissionPolicy {
method == "getRoutingExplanation" ||
method == "getReviewPolicy" ||
method == "getBlockers" ||
method == "getProgress") {
method == "getProgress" ||
method == "getEventStream" ||
method == "getRecentEvents") {
return true;
}

68
editor/src/EventStream.h Normal file
View File

@@ -0,0 +1,68 @@
#pragma once
// Step 387: Orchestrator event stream for polling/subscription.
#include "WorkflowOrchestrator.h"
#include <algorithm>
#include <functional>
#include <vector>
struct StreamEvent {
int version = 0;
OrchestratorEvent event;
json toJson() const {
return {
{"version", version},
{"type", event.type},
{"itemId", event.itemId},
{"detail", event.detail},
{"timestamp", event.timestamp}
};
}
};
class EventStream {
public:
using Callback = std::function<void(const StreamEvent&)>;
void emit(const OrchestratorEvent& event) {
StreamEvent se;
se.version = ++version_;
se.event = event;
events_.push_back(se);
for (const auto& cb : subscribers_) {
cb(se);
}
}
std::vector<StreamEvent> poll(int sinceVersion) const {
std::vector<StreamEvent> out;
for (const auto& se : events_) {
if (se.version > sinceVersion) out.push_back(se);
}
return out;
}
void subscribe(Callback cb) {
subscribers_.push_back(std::move(cb));
}
int getVersion() const { return version_; }
std::vector<StreamEvent> getRecent(int count) const {
if (count <= 0) return {};
int n = std::min(count, static_cast<int>(events_.size()));
return std::vector<StreamEvent>(events_.end() - n, events_.end());
}
static json toJson(const std::vector<StreamEvent>& events) {
json arr = json::array();
for (const auto& e : events) arr.push_back(e.toJson());
return arr;
}
private:
int version_ = 0;
std::vector<StreamEvent> events_;
std::vector<Callback> subscribers_;
};

View File

@@ -2262,6 +2262,12 @@ inline json handleHeadlessAgentRequest(HeadlessEditorState& state,
std::string bufferId = state.activeBuffer->path;
int count = state.workflow->populateFromSkeleton(state.activeAST(), bufferId);
state.workflowProgress = WorkflowProgress(count);
state.eventStream.emit({
"workflow.created",
"",
{{"projectName", projectName}, {"itemCount", count}},
workItemTimestamp()
});
auto stats = state.workflow->getStats();
return headlessRpcResult(id, {

View File

@@ -44,6 +44,7 @@
#include "ContextAssembler.h"
#include "ReviewGate.h"
#include "WorkflowProgress.h"
#include "EventStream.h"
#include <nlohmann/json.hpp>
#include <string>
@@ -143,6 +144,7 @@ struct HeadlessEditorState {
bool verbose = false;
std::optional<WorkflowState> workflow;
std::optional<WorkflowProgress> workflowProgress;
EventStream eventStream;
RoutingEngine routingEngine;
WorkerRegistry workerRegistry = WorkerRegistry::getDefaultRegistry();
ContextAssembler contextAssembler;

View File

@@ -43,6 +43,28 @@ static inline json orchestratorBlockersToJson(const std::vector<BlockerInfo>& bl
return arr;
}
static inline std::string streamTypeForOrchestratorType(const std::string& type) {
if (type == "routed") return "task.routed";
if (type == "context-assembled") return "task.context-assembled";
if (type == "executed") return "task.executed";
if (type == "auto-approved") return "task.auto-approved";
if (type == "sent-to-review") return "task.sent-to-review";
if (type == "completed") return "task.completed";
if (type == "rejected") return "task.rejected";
if (type == "escalated") return "task.escalated";
if (type == "blocked") return "workflow.blocked";
return type;
}
static inline void emitWorkflowEvents(HeadlessEditorState& state,
const std::vector<OrchestratorEvent>& events) {
for (const auto& event : events) {
OrchestratorEvent mapped = event;
mapped.type = streamTypeForOrchestratorType(event.type);
state.eventStream.emit(mapped);
}
}
static inline std::map<std::string, BufferInfo> collectOrchestratorBufferInfos(
HeadlessEditorState& state) {
std::map<std::string, BufferInfo> infos;
@@ -81,6 +103,7 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
state.workflowProgress = WorkflowProgress(state.workflow->getStats().total);
}
state.workflowProgress->recordEvent(event);
emitWorkflowEvents(state, {event});
return orchestratorRpcResult(id, {{"event", orchestratorEventToJson(event)}});
}
@@ -103,6 +126,7 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
for (const auto& event : batch.events) {
state.workflowProgress->recordEvent(event);
}
emitWorkflowEvents(state, batch.events);
return orchestratorRpcResult(id, {
{"events", orchestratorEventsToJson(batch.events)},
@@ -148,6 +172,7 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
for (const auto& event : allEvents) {
state.workflowProgress->recordEvent(event);
}
emitWorkflowEvents(state, allEvents);
return orchestratorRpcResult(id, {
{"events", orchestratorEventsToJson(allEvents)},
@@ -197,6 +222,12 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
orchestrator.setBuffers(collectOrchestratorBufferInfos(state));
snapshot.blockers = orchestrator.getBlockers();
}
state.eventStream.emit({
"workflow.progress",
"",
snapshot.toJson(),
workItemTimestamp()
});
return orchestratorRpcResult(id, snapshot.toJson());
}
@@ -270,6 +301,9 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
events.push_back({"rejected", itemId,
{{"reason", feedback}},
workItemTimestamp()});
events.push_back({"escalated", itemId,
{{"reason", "validation-failed-requeue"}},
workItemTimestamp()});
} else if (acceptance.autoApproved) {
state.workflow->queue.complete(itemId);
events.push_back({"auto-approved", itemId,
@@ -295,6 +329,13 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
for (const auto& event : events) {
state.workflowProgress->recordEvent(event);
}
emitWorkflowEvents(state, events);
if (state.workflow->getStats().total > 0 &&
state.workflow->getStats().complete == state.workflow->getStats().total) {
state.eventStream.emit({"workflow.complete", "", json::object(),
workItemTimestamp()});
}
auto latest = state.workflow->queue.getItem(itemId);
return orchestratorRpcResult(id, {
@@ -305,5 +346,29 @@ inline std::optional<json> tryHandleHeadlessOrchestratorRPC(
});
}
if (method == "getEventStream") {
if (!AgentPermissionPolicy::canInvoke(role, method))
return orchestratorRpcError(id, -32031, "Role not permitted");
auto params = request.contains("params") ? request["params"] : json::object();
int sinceVersion = params.value("sinceVersion", 0);
auto events = state.eventStream.poll(sinceVersion);
return orchestratorRpcResult(id, {
{"version", state.eventStream.getVersion()},
{"events", EventStream::toJson(events)}
});
}
if (method == "getRecentEvents") {
if (!AgentPermissionPolicy::canInvoke(role, method))
return orchestratorRpcError(id, -32031, "Role not permitted");
auto params = request.contains("params") ? request["params"] : json::object();
int count = params.value("count", 20);
auto events = state.eventStream.getRecent(count);
return orchestratorRpcResult(id, {
{"version", state.eventStream.getVersion()},
{"events", EventStream::toJson(events)}
});
}
return std::nullopt;
}

View File

@@ -1600,6 +1600,33 @@ private:
[this](const json& args) {
return callWhetstone("submitExternalResult", args);
};
// whetstone_get_event_stream
tools_.push_back({"whetstone_get_event_stream",
"Poll workflow events emitted since a stream version. "
"Returns normalized event records for visualization and clients.",
{{"type", "object"}, {"properties", {
{"sinceVersion", {{"type", "integer"},
{"description", "Return events with version > sinceVersion"}}}
}}}
});
toolHandlers_["whetstone_get_event_stream"] =
[this](const json& args) {
return callWhetstone("getEventStream", args);
};
// whetstone_get_recent_events
tools_.push_back({"whetstone_get_recent_events",
"Get the latest N workflow events from the event stream.",
{{"type", "object"}, {"properties", {
{"count", {{"type", "integer"},
{"description", "How many most-recent events to return (default 20)"}}}
}}}
});
toolHandlers_["whetstone_get_recent_events"] =
[this](const json& args) {
return callWhetstone("getRecentEvents", args);
};
}
void registerReviewTools() {

View File

@@ -0,0 +1,204 @@
// Step 387: Orchestrator event stream (12 tests)
#include <cassert>
#include <iostream>
#include <string>
#include "EventStream.h"
#include "HeadlessEditorState.h"
#include "HeadlessAgentRPCHandler.h"
#include "MCPServer.h"
static WorkItem makeItem(const std::string& id,
const std::string& workerType = "template",
const std::string& contextWidth = "local") {
WorkItem item;
item.id = id;
item.nodeId = id + "_node";
item.nodeName = "getValue";
item.nodeType = "Function";
item.bufferId = "main.py";
item.workerType = workerType;
item.contextWidth = contextWidth;
item.priority = "medium";
item.status = WI_PENDING;
item.createdAt = workItemTimestamp();
return item;
}
static HeadlessEditorState makeState() {
HeadlessEditorState state;
state.defaultLanguage = "python";
state.openBuffer("main.py", "def a():\n return 1\n", "python");
state.setAgentRole("test-session", AgentRole::Generator);
state.workflow = WorkflowState("events");
return state;
}
static json rpc(HeadlessEditorState& state, const std::string& method,
const json& params = json::object()) {
json request = {{"jsonrpc", "2.0"}, {"id", 1}, {"method", method},
{"params", params}};
return handleHeadlessAgentRequest(state, request, "test-session");
}
int main() {
int passed = 0;
// Test 1: emit + poll returns emitted event
{
EventStream stream;
stream.emit({"task.executed", "a", json::object(), workItemTimestamp()});
auto events = stream.poll(0);
assert(events.size() == 1);
assert(events[0].event.type == "task.executed");
std::cout << "Test 1 PASSED: emit/poll basic behavior\n";
passed++;
}
// Test 2: poll since version returns only new events
{
EventStream stream;
stream.emit({"x", "a", json::object(), workItemTimestamp()}); // v1
stream.emit({"y", "b", json::object(), workItemTimestamp()}); // v2
auto newer = stream.poll(1);
assert(newer.size() == 1);
assert(newer[0].event.type == "y");
std::cout << "Test 2 PASSED: poll filters by version\n";
passed++;
}
// Test 3: stream version increments with emits
{
EventStream stream;
assert(stream.getVersion() == 0);
stream.emit({"a", "", json::object(), workItemTimestamp()});
stream.emit({"b", "", json::object(), workItemTimestamp()});
assert(stream.getVersion() == 2);
std::cout << "Test 3 PASSED: version tracking\n";
passed++;
}
// Test 4: subscribe callback is invoked
{
EventStream stream;
int callbackCount = 0;
stream.subscribe([&callbackCount](const StreamEvent&) {
callbackCount++;
});
stream.emit({"a", "", json::object(), workItemTimestamp()});
assert(callbackCount == 1);
std::cout << "Test 4 PASSED: subscribe callback fires\n";
passed++;
}
// Test 5: orchestration emits normalized task event types
{
auto state = makeState();
state.workflow->queue.enqueue(makeItem("t1", "template"));
rpc(state, "orchestrateAdvance");
auto resp = rpc(state, "getEventStream", {{"sinceVersion", 0}});
assert(resp.contains("result"));
bool hasRouted = false, hasExecuted = false, hasCompleted = false;
for (const auto& e : resp["result"]["events"]) {
std::string t = e.value("type", "");
if (t == "task.routed") hasRouted = true;
if (t == "task.executed") hasExecuted = true;
if (t == "task.completed") hasCompleted = true;
}
assert(hasRouted && hasExecuted && hasCompleted);
std::cout << "Test 5 PASSED: normalized task event types emitted\n";
passed++;
}
// Test 6: getRecent returns requested event count
{
EventStream stream;
stream.emit({"e1", "", json::object(), workItemTimestamp()});
stream.emit({"e2", "", json::object(), workItemTimestamp()});
stream.emit({"e3", "", json::object(), workItemTimestamp()});
auto recent = stream.getRecent(2);
assert(recent.size() == 2);
assert(recent[0].event.type == "e2");
assert(recent[1].event.type == "e3");
std::cout << "Test 6 PASSED: getRecent count + ordering\n";
passed++;
}
// Test 7: empty stream poll returns no events
{
EventStream stream;
auto events = stream.poll(0);
assert(events.empty());
std::cout << "Test 7 PASSED: empty stream poll behavior\n";
passed++;
}
// Test 8: high-frequency polling does not duplicate results
{
EventStream stream;
stream.emit({"a", "", json::object(), workItemTimestamp()});
auto first = stream.poll(0);
auto second = stream.poll(stream.getVersion());
assert(first.size() == 1);
assert(second.empty());
std::cout << "Test 8 PASSED: no duplicate events after poll version update\n";
passed++;
}
// Test 9: stream JSON serialization includes required fields
{
EventStream stream;
stream.emit({"task.executed", "x", json{{"k", "v"}}, workItemTimestamp()});
auto j = EventStream::toJson(stream.poll(0));
assert(j.is_array());
assert(j[0].contains("version"));
assert(j[0].contains("type"));
assert(j[0].contains("timestamp"));
std::cout << "Test 9 PASSED: stream JSON serialization\n";
passed++;
}
// Test 10: MCP tool registration includes event stream tools
{
MCPServer server;
bool hasStream = false, hasRecent = false;
for (const auto& t : server.getTools()) {
if (t.name == "whetstone_get_event_stream") hasStream = true;
if (t.name == "whetstone_get_recent_events") hasRecent = true;
}
assert(hasStream && hasRecent);
std::cout << "Test 10 PASSED: MCP event stream tools registered\n";
passed++;
}
// Test 11: getEventStream RPC returns only events since version
{
auto state = makeState();
state.workflow->queue.enqueue(makeItem("t1", "template"));
rpc(state, "orchestrateAdvance");
auto first = rpc(state, "getEventStream", {{"sinceVersion", 0}});
int v = first["result"]["version"].get<int>();
auto second = rpc(state, "getEventStream", {{"sinceVersion", v}});
assert(first["result"]["events"].size() > 0);
assert(second["result"]["events"].empty());
std::cout << "Test 11 PASSED: getEventStream since-version behavior\n";
passed++;
}
// Test 12: getRecentEvents RPC honors count parameter
{
auto state = makeState();
state.workflow->queue.enqueue(makeItem("a", "template"));
state.workflow->queue.enqueue(makeItem("b", "template"));
rpc(state, "orchestrateAdvance");
auto recent = rpc(state, "getRecentEvents", {{"count", 2}});
assert(recent.contains("result"));
assert(recent["result"]["events"].size() <= 2);
std::cout << "Test 12 PASSED: getRecentEvents count behavior\n";
passed++;
}
std::cout << "\nResults: " << passed << "/12\n";
assert(passed == 12);
return 0;
}

View File

@@ -3146,6 +3146,60 @@ validation failure.
- `editor/src/HeadlessOrchestratorRPC.h` within header-size limit (`309` <= `600`)
- `editor/tests/step386_test.cpp` within test-file size guidance (`269` lines)
### Step 387: Orchestrator Event Stream
**Status:** PASS (12/12 tests)
Added a versioned event stream for orchestration telemetry and exposed polling
APIs through RPC + MCP for client-side visualization and progress UIs.
**Files created:**
- `editor/src/EventStream.h` — versioned stream:
- `emit(OrchestratorEvent)`
- `poll(sinceVersion)`
- `subscribe(callback)`
- `getVersion()`
- `getRecent(count)`
- stream-event JSON serialization
- `editor/tests/step387_test.cpp` — 12 tests covering:
1. emit/poll baseline
2. since-version filtering
3. version tracking
4. subscriber callback firing
5. normalized orchestration event-type emission
6. recent-event ordering/count
7. empty stream behavior
8. no-duplicate polling under frequent reads
9. event JSON structure
10. MCP event-stream tool registration
11. `getEventStream` since-version behavior
12. `getRecentEvents` count behavior
**Files modified:**
- `editor/src/HeadlessEditorState.h` — add `EventStream eventStream`
- `editor/src/HeadlessAgentRPCHandler.h` — emit `workflow.created` event on workflow init
- `editor/src/HeadlessOrchestratorRPC.h` — event emission mapping and new RPC methods:
- `getEventStream`
- `getRecentEvents`
- mapped event names (`task.*`, `workflow.*`) and `workflow.progress/complete` emission
- `editor/src/AgentPermissionPolicy.h` — read permissions for
`getEventStream` / `getRecentEvents`
- `editor/src/MCPServer.h` — MCP tools:
- `whetstone_get_event_stream`
- `whetstone_get_recent_events`
- `editor/CMakeLists.txt``step387_test` target
**Verification run:**
- `step387_test` — PASS (12/12) new step coverage
- `step386_test` — PASS (12/12) regression coverage
- `step382_test` — PASS (12/12) regression coverage
**Architecture gate check:**
- `editor/src/EventStream.h` within header-size limit (`68` <= `600`)
- `editor/src/HeadlessOrchestratorRPC.h` within header-size limit (`374` <= `600`)
- `editor/tests/step387_test.cpp` within test-file size guidance (`204` lines)
- Legacy oversized header persists:
- `editor/src/MCPServer.h` (`1679` > `600`)
# Roadmap Planning — Sprints 12-25+
## Status: Planning Complete (Sprints 12-19 detailed, 20-25 in roadmap.md)