Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
44 commits
Select commit Hold shift + click to select a range
e0f65e2
feat: add hook execution use case
engkimo Jun 29, 2026
640a031
feat: wire chat hook runner
engkimo Jun 29, 2026
d0bf5cb
feat: add shell hook executor
engkimo Jun 30, 2026
2daa5a5
feat: add chat hook execution mode
engkimo Jun 30, 2026
1124e76
feat: add manual hook run cli
engkimo Jun 30, 2026
83d6327
feat: add repl hook run command
engkimo Jun 30, 2026
a095909
feat: add repl tool run noop
engkimo Jul 1, 2026
9ac5ee5
feat: add repl tool laee opt in
engkimo Jul 1, 2026
c270f2f
feat: report repl tool failures
engkimo Jul 1, 2026
ec044f2
feat: add chat permission mode option
engkimo Jul 7, 2026
bb1612b
feat: stream and resume scoped Codex sessions
engkimo Jul 15, 2026
0acdfeb
fix: clean up cancelled native processes
engkimo Jul 15, 2026
c66496c
feat: scope Claude Code execution
engkimo Jul 15, 2026
7df26a3
feat: stream and resume Claude sessions
engkimo Jul 15, 2026
c00ad65
fix: pin native resume to provider
engkimo Jul 15, 2026
61aff96
Record cancelled native turns
engkimo Jul 15, 2026
2ddaf60
Cancel active chat turns in place
engkimo Jul 16, 2026
ef5bf18
Add authenticated chat turn control
engkimo Jul 16, 2026
2359b36
Steer active native chat sessions
engkimo Jul 17, 2026
4e0c66c
Compare recorded agent CLI trials
engkimo Jul 18, 2026
88028ac
Record isolated agent CLI trials
engkimo Jul 18, 2026
9eba8cd
Adjudicate recorded agent CLI trials
engkimo Jul 20, 2026
5678d43
Rehearse agent CLI receipts locally
engkimo Jul 20, 2026
88326ae
Record Phase 43 publication checkpoint
engkimo Jul 21, 2026
d88f1f0
Preflight agent CLI benchmark campaigns
engkimo Jul 21, 2026
a14f814
Govern agent CLI benchmark reviews
engkimo Jul 21, 2026
254a05d
Verify signed benchmark reviews
engkimo Jul 21, 2026
1a68e83
Anchor benchmark campaign provenance
engkimo Jul 21, 2026
60ca2bf
Add authority root transparency
engkimo Jul 21, 2026
21205ba
Add witnessed consistency checkpoints
engkimo Jul 21, 2026
e52567e
Add durable checkpoint registry
engkimo Jul 22, 2026
5aa2cc2
Add authenticated checkpoint range sync
engkimo Jul 22, 2026
48f4a67
Add peer trust rollover continuity
engkimo Jul 22, 2026
f4c7556
Add authenticated checkpoint gossip
engkimo Jul 22, 2026
3e97c7f
Add resumable checkpoint catch-up
engkimo Jul 22, 2026
e324814
Add peer-signed mutual TLS gossip
engkimo Jul 23, 2026
f414a92
Route checkpoint sync over mutual TLS
engkimo Jul 23, 2026
3d1e81d
Normalize colored CLI help tests
engkimo Jul 23, 2026
628aa28
Add peer-signed TLS revocation policy
engkimo Jul 23, 2026
dd61ce2
Expose TLS revocations in trust CLI
engkimo Jul 24, 2026
687446d
Add TLS revocation issuance API
engkimo Jul 25, 2026
e77b5f6
Add TLS revocation CLI workflow
engkimo Jul 25, 2026
d986e03
Cover TLS revocation CLI round trip
engkimo Jul 25, 2026
fedecdc
Expose TLS expiry warnings in status
engkimo Jul 25, 2026
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
815 changes: 815 additions & 0 deletions .serena/memories/morphic_control_plane_strategy.md

Large diffs are not rendered by default.

112 changes: 112 additions & 0 deletions application/use_cases/execute_chat_hook.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
"""Execute validated chat hooks through a hook executor port."""

from __future__ import annotations

from dataclasses import dataclass

from domain.entities.chat_event import ChatEvent, ChatEventType
from domain.entities.chat_session import ChatSession
from domain.entities.hook import (
HookDefinition,
HookDiagnostic,
HookExecutionRequest,
HookExecutionResult,
HookType,
)
from domain.ports.chat_session_store import ChatSessionStorePort
from domain.ports.hook_executor import HookExecutorPort
from domain.ports.hook_registry import HookRegistryPort


@dataclass(frozen=True)
class ExecuteChatHookResult:
session: ChatSession
events: list[ChatEvent]
diagnostics: list[HookDiagnostic]
hook_results: list[HookExecutionResult]


class ExecuteChatHookUseCase:
"""Execute hooks already accepted by workspace validation policy."""

def __init__(
self,
*,
session_store: ChatSessionStorePort,
hook_registry: HookRegistryPort,
hook_executor: HookExecutorPort,
) -> None:
self._session_store = session_store
self._hook_registry = hook_registry
self._hook_executor = hook_executor

async def execute(
self,
*,
session: ChatSession,
hook_type: HookType,
) -> ExecuteChatHookResult:
diagnostics = self._hook_registry.validate()
failures = [diagnostic for diagnostic in diagnostics if diagnostic.status == "FAIL"]
if failures:
names = ", ".join(diagnostic.name for diagnostic in failures)
raise ValueError(f"Hook diagnostics failed: {names}")

current = session
events: list[ChatEvent] = []
hook_results: list[HookExecutionResult] = []

for hook in self._hook_registry.hooks_for(hook_type):
if not hook.enabled:
current, skipped_event = current.record_event(
ChatEventType.HOOK_EXECUTION_SKIPPED,
self._skipped_payload_for(hook),
)
events.append(skipped_event)
continue

request = HookExecutionRequest(
session_id=current.id,
hook_name=hook.name,
hook_type=hook.hook_type,
command=hook.command,
source_path=hook.source_path,
)
current, requested_event = current.record_event(
ChatEventType.HOOK_EXECUTION_REQUESTED,
request.model_dump(mode="json"),
)
events.append(requested_event)

hook_result = await self._hook_executor.execute(request)
hook_results.append(hook_result)
current, completed_event = current.record_event(
ChatEventType.HOOK_EXECUTION_COMPLETED,
{
**hook_result.model_dump(mode="json"),
"hook_name": hook.name,
"hook_type": hook.hook_type.value,
"source_path": hook.source_path,
},
)
events.append(completed_event)

for event in events:
await self._session_store.append_event(event)

return ExecuteChatHookResult(
session=current,
events=events,
diagnostics=diagnostics,
hook_results=hook_results,
)

def _skipped_payload_for(self, hook: HookDefinition) -> dict[str, str]:
return {
"command": hook.command,
"hook_name": hook.name,
"hook_type": hook.hook_type.value,
"reason": "disabled",
"source_path": hook.source_path,
"status": "skipped",
}
36 changes: 33 additions & 3 deletions application/use_cases/execute_chat_tool.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,11 @@
from dataclasses import dataclass
from typing import Any

from application.use_cases.execute_chat_hook import ExecuteChatHookUseCase
from application.use_cases.plan_chat_hooks import PlanChatHooksUseCase
from domain.entities.chat_event import ChatEvent, ChatEventType
from domain.entities.chat_session import ChatSession, PermissionMode
from domain.entities.hook import HookType
from domain.entities.hook import HookExecutionResult, HookType
from domain.ports.chat_session_store import ChatSessionStorePort
from domain.ports.tool_executor import (
ToolExecutionRequest,
Expand All @@ -35,10 +36,12 @@ def __init__(
session_store: ChatSessionStorePort,
tool_executor: ToolExecutorPort,
hook_planner: PlanChatHooksUseCase | None = None,
hook_runner: ExecuteChatHookUseCase | None = None,
) -> None:
self._session_store = session_store
self._tool_executor = tool_executor
self._hook_planner = hook_planner
self._hook_runner = hook_runner

async def execute(
self,
Expand All @@ -57,7 +60,15 @@ async def execute(

current = session
events: list[ChatEvent] = []
if self._hook_planner is not None:
if self._hook_runner is not None:
hook_result = await self._hook_runner.execute(
session=current,
hook_type=HookType.PRE_TOOL,
)
current = hook_result.session
events.extend(hook_result.events)
self._raise_if_hook_failed(hook_result.hook_results, HookType.PRE_TOOL)
elif self._hook_planner is not None:
hook_result = await self._hook_planner.execute(
session=current,
hook_type=HookType.PRE_TOOL,
Expand Down Expand Up @@ -105,7 +116,17 @@ async def execute(
)
tool_events.append(verification_event)

if self._hook_planner is not None:
if self._hook_runner is not None:
for event in tool_events:
await self._session_store.append_event(event)
events.extend(tool_events)
hook_result = await self._hook_runner.execute(
session=current,
hook_type=HookType.POST_TOOL,
)
current = hook_result.session
events.extend(hook_result.events)
elif self._hook_planner is not None:
for event in tool_events:
await self._session_store.append_event(event)
events.extend(tool_events)
Expand Down Expand Up @@ -136,3 +157,12 @@ def _blocked_by_read_only(
if session.permission_mode is not PermissionMode.READ_ONLY:
return False
return risk_level > RiskLevel.SAFE or tool_name not in _READ_ONLY_TOOLS

def _raise_if_hook_failed(
self,
hook_results: list[HookExecutionResult],
hook_type: HookType,
) -> None:
failed = [result for result in hook_results if not result.success]
if failed:
raise RuntimeError(f"{hook_type.value} hook failed")
9 changes: 3 additions & 6 deletions application/use_cases/resume_chat_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,15 +40,12 @@ async def execute(self, session_id: str) -> ResumeChatSessionResult:
if isinstance(mode_value, str):
permission_mode = PermissionMode(mode_value)

status = ChatSessionStatus.ACTIVE
if any(event.type is ChatEventType.SESSION_ENDED for event in events):
status = ChatSessionStatus.ENDED

session = ChatSession(
id=resolved_session_id,
goal=goal,
permission_mode=permission_mode,
status=status,
next_sequence=max(event.sequence for event in events) + 1,
status=ChatSessionStatus.ACTIVE,
)
for event in events:
session = session.replay_event(event)
return ResumeChatSessionResult(session=session, events=events)
116 changes: 110 additions & 6 deletions application/use_cases/route_to_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,19 @@
from dataclasses import dataclass

from application.use_cases.run_council_debate import RunCouncilDebateUseCase
from domain.entities.chat_session import PermissionMode
from domain.entities.cognitive import AgentAction, Decision
from domain.entities.council import SubtaskBrief
from domain.ports.agent_affinity_repository import AgentAffinityRepository
from domain.ports.agent_engine import AgentEngineCapabilities, AgentEnginePort, AgentEngineResult
from domain.ports.agent_engine import (
AgentEngineCapabilities,
AgentEngineEventSinkPort,
AgentEnginePort,
AgentEngineResult,
ResumableStreamingScopedAgentEnginePort,
ScopedAgentEnginePort,
StreamingScopedAgentEnginePort,
)
from domain.ports.context_adapter import ContextAdapterPort
from domain.ports.engine_cost_recorder import EngineCostRecorderPort
from domain.ports.shared_task_state_repository import SharedTaskStateRepository
Expand Down Expand Up @@ -108,6 +117,11 @@ async def execute(
timeout_seconds: float = 300.0,
context: str | None = None,
task_id: str | None = None,
workspace_root: str | None = None,
permission_mode: PermissionMode | None = None,
event_sink: AgentEngineEventSinkPort | None = None,
resume_session_id: str | None = None,
resume_engine: AgentEngineType | None = None,
) -> AgentEngineResult:
"""Route to best available engine and execute.

Expand All @@ -122,6 +136,24 @@ async def execute(
BUG-003: Every attempt is recorded as a FallbackAttempt for transparency.
BUG-002: Successful engine costs are recorded via CostTracker.
"""
if resume_session_id is not None and (
event_sink is None
or workspace_root is None
or permission_mode is None
or resume_engine is None
):
raise ValueError(
"native resume requires engine, streaming, workspace, and permission context"
)
if resume_session_id is None and resume_engine is not None:
raise ValueError("resume engine requires a native session id")
if (
resume_engine is not None
and preferred_engine is not None
and resume_engine is not preferred_engine
):
raise ValueError("preferred engine must match native resume engine")

# Extract topic for affinity lookup
topic = TopicExtractor.extract(task)

Expand All @@ -142,6 +174,15 @@ async def execute(
last_result: AgentEngineResult | None = None

for engine_type in chain:
if resume_engine is not None and engine_type is not resume_engine:
attempts.append(
FallbackAttempt(
engine=engine_type.value,
attempted=False,
skip_reason="resume_engine_mismatch",
)
)
continue
driver = self._drivers.get(engine_type)
if driver is None:
logger.debug("Engine %s not registered, skipping", engine_type.value)
Expand Down Expand Up @@ -175,11 +216,74 @@ async def execute(

start = time.monotonic()
try:
result = await driver.run_task(
task=effective_task,
model=model,
timeout_seconds=timeout_seconds,
)
if workspace_root is not None or permission_mode is not None:
if (
not isinstance(driver, ScopedAgentEnginePort)
or workspace_root is None
or permission_mode is None
):
attempts.append(
FallbackAttempt(
engine=engine_type.value,
attempted=False,
skip_reason="scoped_execution_unsupported",
)
)
continue
if event_sink is not None:
if resume_session_id is not None:
if not isinstance(
driver, ResumableStreamingScopedAgentEnginePort
):
attempts.append(
FallbackAttempt(
engine=engine_type.value,
attempted=False,
skip_reason="native_resume_unsupported",
)
)
continue
result = await driver.resume_task_scoped_stream(
task=effective_task,
resume_session_id=resume_session_id,
workspace_root=workspace_root,
permission_mode=permission_mode,
event_sink=event_sink,
model=model,
timeout_seconds=timeout_seconds,
)
elif not isinstance(driver, StreamingScopedAgentEnginePort):
attempts.append(
FallbackAttempt(
engine=engine_type.value,
attempted=False,
skip_reason="streaming_scoped_execution_unsupported",
)
)
continue
else:
result = await driver.run_task_scoped_stream(
task=effective_task,
workspace_root=workspace_root,
permission_mode=permission_mode,
event_sink=event_sink,
model=model,
timeout_seconds=timeout_seconds,
)
else:
result = await driver.run_task_scoped(
task=effective_task,
workspace_root=workspace_root,
permission_mode=permission_mode,
model=model,
timeout_seconds=timeout_seconds,
)
else:
result = await driver.run_task(
task=effective_task,
model=model,
timeout_seconds=timeout_seconds,
)
except Exception as exc:
elapsed = time.monotonic() - start
logger.warning(
Expand Down
Loading
Loading