diff --git a/editor/CMakeLists.txt b/editor/CMakeLists.txt index 3a2b8aa..7233efb 100644 --- a/editor/CMakeLists.txt +++ b/editor/CMakeLists.txt @@ -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) diff --git a/editor/src/AgentPermissionPolicy.h b/editor/src/AgentPermissionPolicy.h index d6614ab..bc921f8 100644 --- a/editor/src/AgentPermissionPolicy.h +++ b/editor/src/AgentPermissionPolicy.h @@ -72,7 +72,9 @@ struct AgentPermissionPolicy { method == "getRoutingExplanation" || method == "getReviewPolicy" || method == "getBlockers" || - method == "getProgress") { + method == "getProgress" || + method == "getEventStream" || + method == "getRecentEvents") { return true; } diff --git a/editor/src/EventStream.h b/editor/src/EventStream.h new file mode 100644 index 0000000..56a0fd8 --- /dev/null +++ b/editor/src/EventStream.h @@ -0,0 +1,68 @@ +#pragma once +// Step 387: Orchestrator event stream for polling/subscription. + +#include "WorkflowOrchestrator.h" +#include +#include +#include + +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 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 poll(int sinceVersion) const { + std::vector 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 getRecent(int count) const { + if (count <= 0) return {}; + int n = std::min(count, static_cast(events_.size())); + return std::vector(events_.end() - n, events_.end()); + } + + static json toJson(const std::vector& events) { + json arr = json::array(); + for (const auto& e : events) arr.push_back(e.toJson()); + return arr; + } + +private: + int version_ = 0; + std::vector events_; + std::vector subscribers_; +}; diff --git a/editor/src/HeadlessAgentRPCHandler.h b/editor/src/HeadlessAgentRPCHandler.h index 0be0b8a..289cc8d 100644 --- a/editor/src/HeadlessAgentRPCHandler.h +++ b/editor/src/HeadlessAgentRPCHandler.h @@ -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, { diff --git a/editor/src/HeadlessEditorState.h b/editor/src/HeadlessEditorState.h index 81833d3..e159b2e 100644 --- a/editor/src/HeadlessEditorState.h +++ b/editor/src/HeadlessEditorState.h @@ -44,6 +44,7 @@ #include "ContextAssembler.h" #include "ReviewGate.h" #include "WorkflowProgress.h" +#include "EventStream.h" #include #include @@ -143,6 +144,7 @@ struct HeadlessEditorState { bool verbose = false; std::optional workflow; std::optional workflowProgress; + EventStream eventStream; RoutingEngine routingEngine; WorkerRegistry workerRegistry = WorkerRegistry::getDefaultRegistry(); ContextAssembler contextAssembler; diff --git a/editor/src/HeadlessOrchestratorRPC.h b/editor/src/HeadlessOrchestratorRPC.h index 3aba931..324ad5c 100644 --- a/editor/src/HeadlessOrchestratorRPC.h +++ b/editor/src/HeadlessOrchestratorRPC.h @@ -43,6 +43,28 @@ static inline json orchestratorBlockersToJson(const std::vector& 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& events) { + for (const auto& event : events) { + OrchestratorEvent mapped = event; + mapped.type = streamTypeForOrchestratorType(event.type); + state.eventStream.emit(mapped); + } +} + static inline std::map collectOrchestratorBufferInfos( HeadlessEditorState& state) { std::map infos; @@ -81,6 +103,7 @@ inline std::optional 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 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 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 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 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 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 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; } diff --git a/editor/src/MCPServer.h b/editor/src/MCPServer.h index 5e1b580..7368bb8 100644 --- a/editor/src/MCPServer.h +++ b/editor/src/MCPServer.h @@ -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() { diff --git a/editor/tests/step387_test.cpp b/editor/tests/step387_test.cpp new file mode 100644 index 0000000..62d91c8 --- /dev/null +++ b/editor/tests/step387_test.cpp @@ -0,0 +1,204 @@ +// Step 387: Orchestrator event stream (12 tests) + +#include +#include +#include +#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(); + 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; +} diff --git a/progress.md b/progress.md index e98f128..3db6aa3 100644 --- a/progress.md +++ b/progress.md @@ -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)