Skip to content
Merged
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
5 changes: 4 additions & 1 deletion .trunk/trunk.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,12 @@ lint:
# K8s security best practices - remaining issues are acceptable for internal Tailscale-only services
# Fixed: container security contexts, health probes, RBAC over-permissions
# Remaining: image tags/digests, readOnlyRootFilesystem, NetworkPolicy, etc.
# spikes/**: local-kind spike workspace pods (imagePullPolicy Never, no probes, no NetworkPolicy);
# never deployed outside the local kind cluster
- linters: [checkov, trivy]
paths:
- k8s/**
- spikes/**
# B104: Binding to 0.0.0.0 is intentional for containerized services
# B608: False positives - SQL uses parameterized queries, column names are from internal code
# B110: Intentional silent failure for optional API features
Expand Down Expand Up @@ -56,7 +59,7 @@ lint:
- shfmt@3.6.0
- taplo@0.10.0
- terrascan@1.19.9
- trivy@0.68.2
- trivy@0.74.0
- trufflehog@3.92.4
- yamllint@1.37.1
actions:
Expand Down
182 changes: 180 additions & 2 deletions backend/src/mainloop/api.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
"""FastAPI application with DBOS durable workflows."""

import logging
from dataclasses import asdict
from datetime import datetime
from typing import Any

Expand All @@ -15,6 +16,7 @@
ConversationListResponse,
ConversationResponse,
)
from mainloop.runtime.agent_api import router as agent_api_router
from mainloop.services.chat_handler import process_message
from mainloop.services.github_pr import (
CommitSummary,
Expand All @@ -40,6 +42,7 @@

from models import (
MainThread,
NativeSessionInfo,
Project,
QueueItem,
QueueItemResponse,
Expand Down Expand Up @@ -73,8 +76,7 @@ def _apply_mock_github():
if not settings.use_mock_github:
return

import mainloop.services.github_pr as github_pr
from mainloop.services import github_mock
from mainloop.services import github_mock, github_pr

# Replace functions with mocks
funcs_to_mock = [
Expand Down Expand Up @@ -114,6 +116,15 @@ async def startup_event():
# Launch DBOS
DBOS.launch()

if settings.main_thread_mode == "native":
import asyncio

from mainloop.runtime import native_sessions

app.state.native_reconcile = asyncio.create_task(
native_sessions.reconcile_loop()
)


@app.on_event("shutdown")
async def shutdown_event():
Expand Down Expand Up @@ -216,6 +227,9 @@ async def chat(
if not user_id:
user_id = get_user_id_from_cf_header()

if settings.main_thread_mode == "native":
return await _chat_native(request, user_id)

# Ensure main thread is running (for background coordination)
main_thread_id = get_or_start_main_thread(user_id)

Expand Down Expand Up @@ -290,6 +304,96 @@ async def chat(
)


async def _chat_native(request: ChatRequest, user_id: str) -> ChatResponse:
"""Native main thread: record + deliver to the Claude session under Herdr (ledgered). The
reply is mirrored from the native journal, so the client polls the conversation."""
from mainloop.runtime import delegation, native_sessions

binding = await delegation.ensure_main_session(user_id)
session = await db.get_session(binding["session_id"])
try:
message_id = await native_sessions.submit_message(
binding["session_id"], request.message
)
except ValueError as exc:
raise HTTPException(status_code=409, detail=str(exc)) from exc
return ChatResponse(
conversation_id=session.conversation_id,
pending=True,
delivery_message_id=message_id,
)


class MainThreadInfo(BaseModel):
mode: str
session_id: str | None = None
conversation_id: str | None = None
native: NativeSessionInfo | None = None
topics: list[dict] = []


@app.get("/main-thread", response_model=MainThreadInfo)
async def get_main_thread_info(user_id: str = Header(alias="X-User-ID", default=None)):
"""Main-thread mode, native identity strip, and the topic index."""
if not user_id:
user_id = get_user_id_from_cf_header()
if settings.main_thread_mode != "native":
return MainThreadInfo(mode=settings.main_thread_mode)
from mainloop.runtime import delegation, native_sessions

binding = await delegation.ensure_main_session(user_id)
await native_sessions.sync(binding["session_id"])
session = await db.get_session(binding["session_id"])
topics = await delegation._topic_lines(user_id)
return MainThreadInfo(
mode="native",
session_id=binding["session_id"],
conversation_id=session.conversation_id,
native=await native_sessions.identity(binding["session_id"]),
topics=[asdict(t) for t in topics],
)


@app.post("/main-thread/rotate")
async def rotate_main_thread(user_id: str = Header(alias="X-User-ID", default=None)):
"""Force a rotation now (same path as the automatic trigger); used to prove the cut."""
if not user_id:
user_id = get_user_id_from_cf_header()
from mainloop.runtime import delegation, native_sessions

binding = await delegation.ensure_main_session(user_id)
return await native_sessions.rotate(binding["session_id"], "manual")


@app.get("/topics")
async def list_topics(user_id: str = Header(alias="X-User-ID", default=None)):
"""Topic index with records (notes, decisions, pending intent, reports) for the UI."""
if not user_id:
user_id = get_user_id_from_cf_header()
async with db.connection() as conn:
topics = await conn.fetch(
"SELECT * FROM topics WHERE user_id=$1 ORDER BY updated_at DESC", user_id
)
out = []
for t in topics:
recs = await conn.fetch(
"SELECT id, kind, text, status, session_id, created_at FROM topic_records WHERE topic_id=$1 ORDER BY created_at DESC LIMIT 50",
t["id"],
)
out.append(
{
"id": t["id"],
"name": t["name"],
"status_line": t["status_line"],
"records": [dict(r) for r in recs],
}
)
return out


app.include_router(agent_api_router)


# ============= Conversation Endpoints =============


Expand All @@ -315,6 +419,20 @@ async def get_conversation(conversation_id: str):
if not conversation:
raise HTTPException(status_code=404, detail="Conversation not found")

if settings.main_thread_mode == "native":
from mainloop.runtime import native_sessions

async with db.connection() as conn:
main_sid = await conn.fetchval(
"""SELECT b.session_id FROM native_bindings b JOIN sessions s ON s.id=b.session_id
WHERE b.role='main' AND s.conversation_id=$1""",
conversation_id,
)
if main_sid:
await native_sessions.sync(
main_sid
) # mirror new native-journal evidence first

messages = await db.get_messages(conversation_id)
return ConversationResponse(
conversation=conversation,
Expand Down Expand Up @@ -564,6 +682,18 @@ async def list_sessions(

session_status = SessionStatus(status) if status else None
sessions = await db.list_sessions(user_id=user_id, status=session_status)
if sessions:
async with db.connection() as conn:
rows = await conn.fetch(
"""SELECT b.session_id, b.parent_session_id, t.name AS topic FROM native_bindings b
LEFT JOIN topics t ON t.id=b.topic_id WHERE b.session_id = ANY($1)""",
[s.id for s in sessions],
)
info = {r["session_id"]: r for r in rows}
for s in sessions:
if s.id in info:
s.parent_session_id = info[s.id]["parent_session_id"]
s.topic = info[s.id]["topic"]
return sessions


Expand Down Expand Up @@ -627,6 +757,15 @@ async def create_session(
)
session = await db.create_session(session)

if request.agent_kind:
# Real native agent under Herdr in the workspace pod (no DBOS worker / K8s Job).
from mainloop.runtime import native_sessions

await native_sessions.create_binding(session.id, request.agent_kind)
await db.update_session(session.id, status=SessionStatus.ACTIVE)
await native_sessions.submit_message(session.id, request.prompt)
return await db.get_session(session.id)

# Start session worker workflow
with SetWorkflowID(session.id):
worker_queue.enqueue(session_worker_workflow, session.id)
Expand Down Expand Up @@ -659,10 +798,38 @@ async def get_session_conversation(session_id: str):
if not session:
raise HTTPException(status_code=404, detail="Session not found")

from mainloop.runtime import native_sessions

if await native_sessions.get_binding(session_id):
await native_sessions.sync(
session_id
) # mirror new native-journal evidence first
session = await db.get_session(session_id)

messages = await db.get_messages(session.conversation_id)
return SessionConversationResponse(session=session, messages=messages)


@app.get("/sessions/{session_id}/native", response_model=NativeSessionInfo)
async def get_session_native(
session_id: str, user_id: str = Header(alias="X-User-ID", default=None)
):
"""Identity strip for a session bound to a native agent under Herdr."""
if not user_id:
user_id = get_user_id_from_cf_header()
owner = await db.get_session(session_id)
if owner is not None and owner.user_id != user_id:
raise HTTPException(status_code=403, detail="Not your session")
from mainloop.runtime import native_sessions

info = await native_sessions.identity(session_id)
if info is None:
raise HTTPException(
status_code=404, detail="Session has no native agent binding"
)
return info


class SessionMessageRequest(BaseModel):
"""Request to send a message to a session."""

Expand All @@ -688,6 +855,17 @@ async def send_session_message(
if session.user_id != user_id:
raise HTTPException(status_code=403, detail="Not your session")

from mainloop.runtime import native_sessions

if await native_sessions.get_binding(session_id):
try:
message_id = await native_sessions.submit_message(
session_id, request.message
)
except ValueError as exc:
raise HTTPException(status_code=409, detail=str(exc)) from exc
return {"status": "ok", "message_id": message_id}

# Save message directly to database (don't rely on workflow)
message = await db.create_message(
conversation_id=session.conversation_id,
Expand Down
21 changes: 21 additions & 0 deletions backend/src/mainloop/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,27 @@ def database_url(self) -> str:
claude_model: str = "sonnet" # Main thread model
claude_worker_model: str = "opus" # Worker model (for background tasks)

# Native agents under Herdr (workspace pod reached over Kubernetes pod-exec)
workspace_namespace: str = "herdr-spike"
workspace_pod: str = "workspace-0"
main_pod: str = (
"main-0" # pod that runs the native main thread (scratch cwd, no repo)
)

# Native main thread (context model). MAIN_THREAD_MODE=native replaces the SDK chat path.
main_thread_mode: str = "sdk" # sdk | native
main_thread_model: str = "sonnet"
main_thread_effort: str = "medium"
# Rotation: cut to a fresh native session when the context grew by this many tokens above
# the lineage's first-turn baseline, or after this many completed turns (whichever first).
main_rotate_tokens: int = 20000
main_rotate_turns: int = 12
main_carry_over_messages: int = 6
native_child_kinds: str = "claude,codex"
agent_token_key: str = (
"" # HMAC key for per-binding agent tokens (falls back to DB password)
)

# GitHub
github_token: str = ""

Expand Down
Loading
Loading