This repository has been archived on 2026-07-21. You can view files and clone it, but cannot push or open issues or pull requests.
ollie-dbus/main.go

1376 lines
34 KiB
Go

package main
import (
"context"
"encoding/json"
"fmt"
"hash/crc32"
"os"
"os/signal"
"path/filepath"
"strings"
"sync"
"syscall"
"time"
"github.com/godbus/dbus/v5"
"github.com/godbus/dbus/v5/introspect"
"ollie/pkg/agent"
"ollie/pkg/backend"
"ollie/pkg/config"
"ollie/pkg/env"
"ollie/pkg/paths"
"ollie/pkg/tools"
"ollie/pkg/tools/execute"
)
const (
busName = "org.ollie.SessionManager"
busPath = "/org/ollie/SessionManager"
busIface = "org.ollie.SessionManager"
)
// managedSession wraps a Core with its context and metadata.
type managedSession struct {
core agent.Core
ctx context.Context
cancel context.CancelFunc
id string
agent string
log []byte
logMu sync.Mutex
peers map[string]bool
}
// chatEmitter batches ChatUpdated D-Bus signals to avoid flooding the bus.
// Text is appended to the session log immediately, but the D-Bus signal is
// coalesced and emitted at most every 50ms.
type chatEmitter struct {
conn *dbus.Conn
sess *managedSession
mu sync.Mutex
buf string
offset int64
timer *time.Timer
running bool
}
func newChatEmitter(conn *dbus.Conn, sess *managedSession) *chatEmitter {
return &chatEmitter{conn: conn, sess: sess}
}
func (e *chatEmitter) append(text string) {
// Write to log immediately so GetChat always returns the latest.
e.sess.logMu.Lock()
if e.buf == "" {
e.offset = int64(len(e.sess.log))
}
e.sess.log = append(e.sess.log, []byte(text)...)
e.sess.logMu.Unlock()
e.mu.Lock()
e.buf += text
if !e.running {
e.running = true
e.timer = time.AfterFunc(50*time.Millisecond, e.flush)
}
e.mu.Unlock()
}
func (e *chatEmitter) flush() {
e.mu.Lock()
text := e.buf
offset := e.offset
e.buf = ""
e.running = false
e.mu.Unlock()
if text != "" {
e.conn.Emit(busPath, busIface+".ChatUpdated", e.sess.id, offset, text)
}
}
// flushSync drains any pending buffered text immediately (used on state transitions).
func (e *chatEmitter) flushSync() {
e.mu.Lock()
if e.timer != nil {
e.timer.Stop()
}
text := e.buf
offset := e.offset
e.buf = ""
e.running = false
e.mu.Unlock()
if text != "" {
e.conn.Emit(busPath, busIface+".ChatUpdated", e.sess.id, offset, text)
}
}
// SessionManager is the D-Bus exported object.
type SessionManager struct {
mu sync.RWMutex
sessions map[string]*managedSession
conn *dbus.Conn
}
func NewSessionManager(conn *dbus.Conn) *SessionManager {
env.EnsureDefaults()
return &SessionManager{
sessions: make(map[string]*managedSession),
conn: conn,
}
}
// --- Session lifecycle ---
func (m *SessionManager) CreateSession(cwd, backendName, modelName, agentName string) (string, *dbus.Error) {
if cwd == "" {
cwd, _ = os.Getwd()
}
cwd = paths.ExpandHome(os.ExpandEnv(cwd))
if agentName == "" {
agentName = "default"
if v := os.Getenv("OLLIE_DEFAULT_AGENT"); v != "" {
agentName = v
}
}
sessID := agent.NewSessionID()
// Load agent config
agentsDir := paths.CfgDir() + "/agents"
var cfg *config.Config
cfgPath := agent.AgentConfigPath(agentsDir, agentName)
if f, err := os.Open(cfgPath); err == nil {
cfg, _ = config.Load(f)
f.Close()
}
// Resolve backend
if backendName == "" && cfg != nil && cfg.Backend != "" {
backendName = cfg.Backend
}
be, err := backend.NewWithName(backendName)
if err != nil {
return "", dbus.MakeFailedError(fmt.Errorf("backend: %w", err))
}
// Resolve model
if modelName == "" && cfg != nil && cfg.Model != "" {
modelName = cfg.Model
}
if modelName == "" {
modelName = os.Getenv("OLLIE_MODEL")
}
if modelName != "" {
be.SetModel(modelName)
}
// Sessions dir
sessionsDir := paths.DataDir() + "/sessions"
os.MkdirAll(sessionsDir, 0700)
// Build dispatcher + runtime
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
"execute": execute.Decl(cwd),
})
rt := agent.BuildRuntime(cfg, newDisp(), cwd, []string{"OLLIE_SESSION_ID=" + sessID})
// Wire OnExit hook so detached process exits get injected into agent context.
// We use a pointer indirection because core hasn't been created yet.
var coreRef agent.Core
if srv, ok := rt.Dispatcher.GetServer("execute"); ok {
type exitHooker interface {
SetOnExit(func(int, int))
}
if eh, ok := srv.(exitHooker); ok {
eh.SetOnExit(func(pid, exitCode int) {
if coreRef == nil {
return
}
var output string
type outputGetter interface {
GetDetachedOutput(int) (string, error)
}
if og, ok := srv.(outputGetter); ok {
output, _ = og.GetDetachedOutput(pid)
}
if len(output) > 2048 {
output = "[...truncated...]\n" + output[len(output)-2048:]
}
msg := fmt.Sprintf("pid %d exited with code %d\n%s", pid, exitCode, output)
coreRef.InjectSystemEvent(msg)
m.conn.Emit(busPath, busIface+".ProcessExited", sessID, int32(pid), int32(exitCode))
})
}
}
// Create the agent core
core := agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: agentName,
AgentsDir: agentsDir,
SessionsDir: sessionsDir,
SessionID: sessID,
CWD: cwd,
Runtime: rt,
NewDispatcher: newDisp,
})
coreRef = core
ctx, cancel := context.WithCancel(context.Background())
sess := &managedSession{
core: core,
ctx: ctx,
cancel: cancel,
id: sessID,
agent: agentName,
}
// Subscribe to events for the chat log + signal dispatch
emitter := newChatEmitter(m.conn, sess)
streamingRole := ""
core.Bus().Subscribe("event", func(ev agent.Event) {
var text string
switch ev.Role {
case "assistant", "reasoning":
if streamingRole != ev.Role {
if streamingRole != "" {
text += "\n\n"
}
text += "[" + ev.Role + "]\n"
streamingRole = ev.Role
}
text += ev.Content
default:
if ev.Role == "usage" {
return
}
if streamingRole != "" {
text += "\n\n"
streamingRole = ""
}
switch ev.Role {
case "user":
text += "[user]\n" + ev.Content + "\n\n"
case "call":
text += "[call:" + ev.Name + "]\n" + ev.Content + "\n\n"
case "tool":
text += "[tool:" + ev.Name + "]\n" + ev.Content + "\n\n"
case "error":
text += "[error]\n" + ev.Content + "\n\n"
case "info":
text += "[info] " + ev.Content + "\n\n"
default:
text += "[" + ev.Role + "]\n" + ev.Content + "\n\n"
}
}
if text == "" {
return
}
emitter.append(text)
})
// Watch state changes
go func() {
current := "idle"
for {
next, ok := core.WaitChange(ctx, agent.WatchState, current)
if !ok {
emitter.flushSync()
return
}
current = next
emitter.flushSync()
m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current)
// Send notification when a turn completes
if current == "idle" {
m.sendNotification("Ollie", "Session "+sess.id+" finished")
m.persistSession(sess.id)
}
}
}()
m.mu.Lock()
m.sessions[sessID] = sess
m.mu.Unlock()
m.conn.Emit(busPath, busIface+".SessionCreated", sessID)
return sessID, nil
}
func (m *SessionManager) ListSessions() ([]string, *dbus.Error) {
m.mu.RLock()
defer m.mu.RUnlock()
result := make([]string, 0, len(m.sessions))
for id, sess := range m.sessions {
result = append(result, fmt.Sprintf("%s\t%s\t%s\t%s",
id, sess.core.State(), sess.core.ModelName(), sess.agent))
}
return result, nil
}
func (m *SessionManager) KillSession(sessionID string) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
if !ok {
m.mu.Unlock()
return false, nil
}
delete(m.sessions, sessionID)
m.mu.Unlock()
sess.core.Interrupt(agent.ErrInterrupted)
sess.core.Close()
sess.cancel()
m.removePersistedSession(sessionID)
m.conn.Emit(busPath, busIface+".SessionKilled", sessionID)
return true, nil
}
func (m *SessionManager) RenameSession(sessionID, newName string) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
if !ok {
m.mu.Unlock()
return false, nil
}
delete(m.sessions, sessionID)
sess.id = newName
m.sessions[newName] = sess
m.mu.Unlock()
if err := sess.core.SetSessionID(newName); err != nil {
m.mu.Lock()
delete(m.sessions, newName)
sess.id = sessionID
m.sessions[sessionID] = sess
m.mu.Unlock()
return false, nil
}
m.conn.Emit(busPath, busIface+".SessionRenamed", sessionID, newName)
return true, nil
}
// --- Interaction ---
func (m *SessionManager) Submit(sessionID, prompt string) (bool, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return false, nil
}
go sess.core.Submit(sess.ctx, prompt)
return true, nil
}
func (m *SessionManager) Interrupt(sessionID string) (bool, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return false, nil
}
sess.core.Interrupt(agent.ErrInterrupted)
return true, nil
}
// --- State queries ---
func (m *SessionManager) GetState(sessionID string) (string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", nil
}
return sess.core.State(), nil
}
func (m *SessionManager) GetChat(sessionID string, offset int64) (string, int64, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", offset, nil
}
sess.logMu.Lock()
if offset < 0 {
offset = 0
}
if offset > int64(len(sess.log)) {
offset = int64(len(sess.log))
}
text := string(sess.log[offset:])
newOffset := int64(len(sess.log))
sess.logMu.Unlock()
return text, newOffset, nil
}
func (m *SessionManager) GetUsage(sessionID string) (string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", nil
}
return sess.core.Usage(), nil
}
func (m *SessionManager) GetCost(sessionID string) (string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", nil
}
return sess.core.Cost(), nil
}
// --- Config ---
func (m *SessionManager) GetConfig(sessionID string) (string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", nil
}
var result string
result += "model=" + sess.core.ModelName() + "\n"
result += "backend=" + sess.core.BackendName() + "\n"
result += "agent=" + sess.agent + "\n"
result += "cwd=" + sess.core.CWD() + "\n"
result += "state=" + sess.core.State() + "\n"
return result, nil
}
func (m *SessionManager) SetConfig(sessionID, key, value string) (bool, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return false, nil
}
switch key {
case "model":
sess.core.Submit(sess.ctx, "/model "+value)
case "cwd":
if err := sess.core.SetCWD(value); err != nil {
return false, nil
}
default:
return false, nil
}
return true, nil
}
// --- Context ---
func (m *SessionManager) GetContext(sessionID string) (string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return "", nil
}
msgs := sess.core.Context()
var result string
for _, msg := range msgs {
data, err := json.Marshal(msg)
if err != nil {
continue
}
result += string(data) + "\n"
}
return result, nil
}
// --- Backends/models/agents ---
func (m *SessionManager) ListBackends() ([]string, *dbus.Error) {
return []string{"ollama", "openai", "openrouter", "anthropic", "copilot", "kiro", "gemini"}, nil
}
func (m *SessionManager) ListModels(sessionID string) ([]string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return nil, nil
}
raw := sess.core.ListModels()
models := strings.Split(strings.TrimSpace(raw), "\n")
if len(models) == 1 && models[0] == "" {
return nil, nil
}
return models, nil
}
func (m *SessionManager) ListAgents() ([]string, *dbus.Error) {
var result []string
for _, dir := range agent.AgentsDirs() {
entries, err := os.ReadDir(dir)
if err != nil {
continue
}
for _, e := range entries {
name := e.Name()
if strings.HasSuffix(name, ".json") {
result = append(result, strings.TrimSuffix(name, ".json"))
}
}
}
return result, nil
}
// --- Peers ---
func (m *SessionManager) PeerAdd(sessionID, peerID string) (bool, *dbus.Error) {
m.mu.RLock()
_, ok := m.sessions[sessionID]
_, peerOk := m.sessions[peerID]
m.mu.RUnlock()
if !ok || !peerOk {
return false, nil
}
m.mu.Lock()
sess := m.sessions[sessionID]
if sess.peers == nil {
sess.peers = make(map[string]bool)
}
sess.peers[peerID] = true
m.mu.Unlock()
return true, nil
}
func (m *SessionManager) PeerRemove(sessionID, peerID string) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
if !ok {
m.mu.Unlock()
return false, nil
}
delete(sess.peers, peerID)
m.mu.Unlock()
return true, nil
}
func (m *SessionManager) PeerList(sessionID string) ([]string, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
m.mu.RUnlock()
if !ok {
return nil, nil
}
var result []string
for peerID := range sess.peers {
result = append(result, peerID)
}
return result, nil
}
func (m *SessionManager) PeerSubmit(sessionID, peerID, prompt string) (bool, *dbus.Error) {
m.mu.RLock()
sess, ok := m.sessions[sessionID]
peer, peerOk := m.sessions[peerID]
m.mu.RUnlock()
if !ok || !peerOk {
return false, nil
}
if !sess.peers[peerID] {
return false, nil
}
go peer.core.Submit(peer.ctx, prompt)
return true, nil
}
// --- Detached Process Management ---
func (m *SessionManager) DetachProcess(sessionID string) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
m.mu.Unlock()
if !ok {
return false, dbus.MakeFailedError(fmt.Errorf("session not found: %s", sessionID))
}
if !sess.core.Detach() {
return false, dbus.MakeFailedError(fmt.Errorf("no running process to detach"))
}
m.conn.Emit(busPath, busIface+".ProcessDetached", sessionID)
return true, nil
}
func (m *SessionManager) ListDetached(sessionID string) ([]string, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
m.mu.Unlock()
if !ok {
return nil, dbus.MakeFailedError(fmt.Errorf("session not found: %s", sessionID))
}
infos := sess.core.ListDetached()
out := make([]string, len(infos))
for i, info := range infos {
status := "running"
if info.Exited {
status = fmt.Sprintf("exited(%d)", info.ExitCode)
}
out[i] = fmt.Sprintf("%d\t%s\t%d\t%s", info.PID, info.Command, info.Started, status)
}
return out, nil
}
func (m *SessionManager) SignalDetached(sessionID string, pid int32, signal string) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
m.mu.Unlock()
if !ok {
return false, dbus.MakeFailedError(fmt.Errorf("session not found: %s", sessionID))
}
var sig int
switch strings.ToUpper(signal) {
case "TERM", "SIGTERM", "15":
sig = 15
case "KILL", "SIGKILL", "9":
sig = 9
default:
return false, dbus.MakeFailedError(fmt.Errorf("unsupported signal: %s (use TERM or KILL)", signal))
}
if err := sess.core.SignalDetached(int(pid), sig); err != nil {
return false, dbus.MakeFailedError(err)
}
return true, nil
}
func (m *SessionManager) GetDetachedOutput(sessionID string, pid int32) (string, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
m.mu.Unlock()
if !ok {
return "", dbus.MakeFailedError(fmt.Errorf("session not found: %s", sessionID))
}
output, err := sess.core.GetDetachedOutput(int(pid))
if err != nil {
return "", dbus.MakeFailedError(err)
}
return output, nil
}
func (m *SessionManager) DismissDetached(sessionID string, pid int32) (bool, *dbus.Error) {
m.mu.Lock()
sess, ok := m.sessions[sessionID]
m.mu.Unlock()
if !ok {
return false, dbus.MakeFailedError(fmt.Errorf("session not found: %s", sessionID))
}
return sess.core.DismissDetached(int(pid)), nil
}
// --- Notifications ---
func (m *SessionManager) sendNotification(title, body string) {
obj := m.conn.Object("org.freedesktop.Notifications", "/org/freedesktop/Notifications")
obj.Call("org.freedesktop.Notifications.Notify", 0,
"ollie-dbus", // app_name
uint32(0), // replaces_id
"system-run", // app_icon
title, // summary
body, // body
[]string{}, // actions
map[string]dbus.Variant{}, // hints
int32(5000), // timeout ms
)
}
// --- Session Persistence ---
// activeSessionsDir returns the directory for persisting active sessions across restarts.
func activeSessionsDir() string {
return filepath.Join(paths.DataDir(), "active-sessions")
}
// persistSession saves a single session to the active-sessions directory.
// Called on every idle transition for crash resilience.
func (m *SessionManager) persistSession(id string) {
if strings.HasSuffix(id, "-copilot") {
return
}
dir := activeSessionsDir()
os.MkdirAll(dir, 0700)
m.mu.RLock()
sess, ok := m.sessions[id]
m.mu.RUnlock()
if !ok {
return
}
path := filepath.Join(dir, id+".json")
if err := sess.core.SaveSession(path); err != nil {
fmt.Fprintf(os.Stderr, "checkpoint session %s: %v\n", id, err)
}
}
// removePersistedSession deletes the checkpoint file for a killed session.
func (m *SessionManager) removePersistedSession(id string) {
path := filepath.Join(activeSessionsDir(), id+".json")
os.Remove(path)
}
// rebuildChatLog reconstructs a human-readable chat log from the tail of the
// persisted message list. This gives restored sessions visible history.
func rebuildChatLog(messages []backend.Message) []byte {
// Show up to the last hotTailSize messages (matches the verbatim context window).
const tailSize = 8
start := 0
if len(messages) > tailSize {
start = len(messages) - tailSize
}
tail := messages[start:]
var buf strings.Builder
if start > 0 {
buf.WriteString(fmt.Sprintf("[info] session restored (%d earlier messages compacted)\n\n", start))
}
for _, msg := range tail {
switch msg.Role {
case "system":
continue // don't show system messages in chat log
case "user":
buf.WriteString("[user]\n" + msg.Content + "\n\n")
case "assistant":
if msg.Reasoning != "" {
buf.WriteString("[reasoning]\n" + msg.Reasoning + "\n\n")
}
if msg.Content != "" {
buf.WriteString("[assistant]\n" + msg.Content + "\n\n")
}
for _, tc := range msg.ToolCalls {
buf.WriteString("[call:" + tc.Name + "]\n" + string(tc.Arguments) + "\n\n")
}
case "tool":
name := "result"
if msg.ToolCallID != "" {
name = msg.ToolCallID
}
// Truncate long tool results in the log
content := msg.Content
if len(content) > 2000 {
content = content[:2000] + "\n... (truncated)"
}
buf.WriteString("[tool:" + name + "]\n" + content + "\n\n")
}
}
return []byte(buf.String())
}
// saveAllSessions persists every active session to disk so they can be restored on restart.
func (m *SessionManager) saveAllSessions() {
dir := activeSessionsDir()
os.MkdirAll(dir, 0700)
// Remove stale files from previous run
entries, _ := os.ReadDir(dir)
for _, e := range entries {
os.Remove(filepath.Join(dir, e.Name()))
}
m.mu.RLock()
defer m.mu.RUnlock()
for id, sess := range m.sessions {
// Skip copilot sessions — they are ephemeral
if strings.HasSuffix(id, "-copilot") {
continue
}
path := filepath.Join(dir, id+".json")
if err := sess.core.SaveSession(path); err != nil {
fmt.Fprintf(os.Stderr, "persist session %s: %v\n", id, err)
}
}
}
// restoreAllSessions loads previously-persisted sessions and recreates them.
func (m *SessionManager) restoreAllSessions() {
dir := activeSessionsDir()
entries, err := os.ReadDir(dir)
if err != nil {
return // no saved sessions
}
for _, e := range entries {
if !strings.HasSuffix(e.Name(), ".json") {
continue
}
path := filepath.Join(dir, e.Name())
ps, err := agent.LoadPersistedSession(path)
if err != nil {
fmt.Fprintf(os.Stderr, "restore session %s: %v\n", e.Name(), err)
continue
}
if err := m.restoreSession(ps); err != nil {
fmt.Fprintf(os.Stderr, "restore session %s: %v\n", ps.ID, err)
}
}
// Clean up persisted files after successful restore
for _, e := range entries {
os.Remove(filepath.Join(dir, e.Name()))
}
}
// restoreSession recreates a single session from persisted state.
func (m *SessionManager) restoreSession(ps *agent.PersistedSession) error {
cwd := ps.CWD
if cwd == "" {
cwd, _ = os.Getwd()
}
agentName := ps.Agent
if agentName == "" {
agentName = "default"
}
sessID := ps.ID
// Load agent config
agentsDir := paths.CfgDir() + "/agents"
var cfg *config.Config
cfgPath := agent.AgentConfigPath(agentsDir, agentName)
if f, err := os.Open(cfgPath); err == nil {
cfg, _ = config.Load(f)
f.Close()
}
// Resolve backend
backendName := ps.Backend
if backendName == "" && cfg != nil && cfg.Backend != "" {
backendName = cfg.Backend
}
be, err := backend.NewWithName(backendName)
if err != nil {
return fmt.Errorf("backend: %w", err)
}
// Resolve model
modelName := ps.Model
if modelName == "" && cfg != nil && cfg.Model != "" {
modelName = cfg.Model
}
if modelName == "" {
modelName = os.Getenv("OLLIE_MODEL")
}
if modelName != "" {
be.SetModel(modelName)
}
// Sessions dir
sessionsDir := paths.DataDir() + "/sessions"
os.MkdirAll(sessionsDir, 0700)
// Build dispatcher + runtime
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
"execute": execute.Decl(cwd),
})
rt := agent.BuildRuntime(cfg, newDisp(), cwd, []string{"OLLIE_SESSION_ID=" + sessID})
// Restore session from persisted messages
restoredSession := agent.RestoreSession(ps.Messages)
if ps.TaskState != nil {
restoredSession.TaskState = ps.TaskState
}
// Create the agent core with restored session
core := agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: agentName,
AgentsDir: agentsDir,
SessionsDir: sessionsDir,
SessionID: sessID,
CWD: cwd,
Session: restoredSession,
Runtime: rt,
NewDispatcher: newDisp,
})
ctx, cancel := context.WithCancel(context.Background())
sess := &managedSession{
core: core,
ctx: ctx,
cancel: cancel,
id: sessID,
agent: agentName,
log: rebuildChatLog(ps.Messages),
}
// Subscribe to events for the chat log + signal dispatch
emitter := newChatEmitter(m.conn, sess)
streamingRole := ""
core.Bus().Subscribe("event", func(ev agent.Event) {
var text string
switch ev.Role {
case "assistant", "reasoning":
if streamingRole != ev.Role {
if streamingRole != "" {
text += "\n\n"
}
text += "[" + ev.Role + "]\n"
streamingRole = ev.Role
}
text += ev.Content
default:
if ev.Role == "usage" {
return
}
if streamingRole != "" {
text += "\n\n"
streamingRole = ""
}
switch ev.Role {
case "user":
text += "[user]\n" + ev.Content + "\n\n"
case "call":
text += "[call:" + ev.Name + "]\n" + ev.Content + "\n\n"
case "tool":
text += "[tool:" + ev.Name + "]\n" + ev.Content + "\n\n"
case "error":
text += "[error]\n" + ev.Content + "\n\n"
case "info":
text += "[info] " + ev.Content + "\n\n"
default:
text += "[" + ev.Role + "]\n" + ev.Content + "\n\n"
}
}
if text == "" {
return
}
emitter.append(text)
})
// Watch state changes
go func() {
current := "idle"
for {
next, ok := core.WaitChange(ctx, agent.WatchState, current)
if !ok {
emitter.flushSync()
return
}
current = next
emitter.flushSync()
m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current)
if current == "idle" {
m.sendNotification("Ollie", "Session "+sess.id+" finished")
m.persistSession(sess.id)
}
}
}()
m.mu.Lock()
m.sessions[sessID] = sess
m.mu.Unlock()
m.conn.Emit(busPath, busIface+".SessionCreated", sessID)
fmt.Printf(" restored session: %s [%s/%s]\n", sessID, agentName, modelName)
return nil
}
// --- Shutdown ---
// --- Completion ---
// Complete performs a single-shot code completion using a dedicated copilot session.
// It takes context about the cursor position and returns raw completion text.
func (m *SessionManager) Complete(cwd, filePath, prefix, suffix, extraContext string) (string, *dbus.Error) {
if cwd == "" {
cwd, _ = os.Getwd()
}
cwd = paths.ExpandHome(os.ExpandEnv(cwd))
model := os.Getenv("OLLIE_COMPLETE_MODEL")
backendName := os.Getenv("OLLIE_COMPLETE_BACKEND")
if model == "" || backendName == "" {
return "", dbus.MakeFailedError(fmt.Errorf("OLLIE_COMPLETE_MODEL and OLLIE_COMPLETE_BACKEND must be set"))
}
// Truncate to context budget
const prefixMax = 12000
const suffixMax = 1000
if len(prefix) > prefixMax {
prefix = prefix[len(prefix)-prefixMax:]
}
if len(suffix) > suffixMax {
suffix = suffix[:suffixMax]
}
// Find or create a copilot session keyed by cwd
sessID := m.findOrCreateCopilot(cwd, backendName, model)
if sessID == "" {
return "", dbus.MakeFailedError(fmt.Errorf("failed to create copilot session"))
}
m.mu.RLock()
sess, ok := m.sessions[sessID]
m.mu.RUnlock()
if !ok {
return "", dbus.MakeFailedError(fmt.Errorf("copilot session disappeared"))
}
// Build FIM prompt
fileHint := ""
if filePath != "" {
fileHint = " in " + filePath
}
contextBlock := ""
if extraContext != "" {
contextBlock = "\n" + extraContext
}
prompt := fmt.Sprintf(`Implement the code at the cursor%s. The prefix ends at the point where new code is needed. Write the implementation — do not echo stubs, TODOs, or placeholder returns from the prefix. Output ONLY raw code. No reasoning, no shell commands, no explanations, no markdown fences, no backticks, no preamble. Your entire response must be valid code that can be inserted directly into the file.
%s
<prefix>
%s
</prefix>
<suffix>
%s
</suffix>`, fileHint, contextBlock, prefix, suffix)
// Submit and wait for completion
sess.core.Submit(sess.ctx, "/clear")
sess.core.Submit(sess.ctx, prompt)
// Wait for idle
for {
state := sess.core.State()
if state == "idle" {
break
}
next, ok := sess.core.WaitChange(sess.ctx, agent.WatchState, state)
if !ok {
return "", dbus.MakeFailedError(fmt.Errorf("copilot session cancelled"))
}
if next == "idle" {
break
}
}
result := sess.core.Reply()
// Strip markdown fences and leaked XML tags
result = stripCompletionNoise(result)
// Strip prefix echo
result = stripPrefixEcho(prefix, result)
return result, nil
}
// findOrCreateCopilot finds an existing copilot session for cwd or creates one.
func (m *SessionManager) findOrCreateCopilot(cwd, backendName, modelName string) string {
// Session ID is deterministic based on cwd
sessID := fmt.Sprintf("%d-copilot", crc32Str(cwd))
m.mu.RLock()
_, exists := m.sessions[sessID]
m.mu.RUnlock()
if exists {
return sessID
}
// Create a new copilot session
agentsDir := paths.CfgDir() + "/agents"
var cfg *config.Config
cfgPath := agent.AgentConfigPath(agentsDir, "copilot")
if f, err := os.Open(cfgPath); err == nil {
cfg, _ = config.Load(f)
f.Close()
}
be, err := backend.NewWithName(backendName)
if err != nil {
return ""
}
be.SetModel(modelName)
sessionsDir := paths.DataDir() + "/sessions"
os.MkdirAll(sessionsDir, 0700)
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
"execute": execute.Decl(cwd),
})
rt := agent.BuildRuntime(cfg, newDisp(), cwd, []string{"OLLIE_SESSION_ID=" + sessID})
core := agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: "copilot",
AgentsDir: agentsDir,
SessionsDir: sessionsDir,
SessionID: sessID,
CWD: cwd,
Runtime: rt,
NewDispatcher: newDisp,
})
ctx, cancel := context.WithCancel(context.Background())
sess := &managedSession{
core: core,
ctx: ctx,
cancel: cancel,
id: sessID,
agent: "copilot",
}
m.mu.Lock()
// Double-check after lock
if _, exists := m.sessions[sessID]; exists {
m.mu.Unlock()
core.Close()
cancel()
return sessID
}
m.sessions[sessID] = sess
m.mu.Unlock()
return sessID
}
// stripCompletionNoise removes markdown fences and leaked XML tags.
func stripCompletionNoise(s string) string {
var lines []string
for _, line := range strings.Split(s, "\n") {
trimmed := strings.TrimSpace(line)
// Skip markdown code fences
if strings.HasPrefix(trimmed, "```") {
continue
}
// Skip leaked XML tags
if trimmed == "<prefix>" || trimmed == "</prefix>" ||
trimmed == "<suffix>" || trimmed == "</suffix>" {
continue
}
// Skip info lines
if strings.HasPrefix(trimmed, ":: ") {
continue
}
lines = append(lines, line)
}
return strings.Join(lines, "\n")
}
// stripPrefixEcho removes echoed prefix from the beginning of the result.
func stripPrefixEcho(prefix, result string) string {
tailMax := 200
if len(prefix) < tailMax {
tailMax = len(prefix)
}
for i := tailMax; i > 0; i-- {
tail := prefix[len(prefix)-i:]
if strings.HasPrefix(result, tail) {
return result[len(tail):]
}
}
return result
}
// crc32Str returns a CRC32 checksum of a string as uint32.
func crc32Str(s string) uint32 {
return crc32.ChecksumIEEE([]byte(s))
}
func (m *SessionManager) Shutdown() {
m.saveAllSessions()
m.mu.Lock()
for id, sess := range m.sessions {
sess.core.Interrupt(agent.ErrInterrupted)
sess.core.Close()
sess.cancel()
delete(m.sessions, id)
}
m.mu.Unlock()
}
// --- Introspection XML ---
const introspectXML = `<node>
<interface name="org.ollie.SessionManager">
<method name="CreateSession">
<arg name="cwd" type="s" direction="in"/>
<arg name="backend" type="s" direction="in"/>
<arg name="model" type="s" direction="in"/>
<arg name="agent" type="s" direction="in"/>
<arg name="session_id" type="s" direction="out"/>
</method>
<method name="ListSessions">
<arg name="sessions" type="as" direction="out"/>
</method>
<method name="KillSession">
<arg name="session_id" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="RenameSession">
<arg name="session_id" type="s" direction="in"/>
<arg name="new_name" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="Submit">
<arg name="session_id" type="s" direction="in"/>
<arg name="prompt" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="Interrupt">
<arg name="session_id" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="GetState">
<arg name="session_id" type="s" direction="in"/>
<arg name="state" type="s" direction="out"/>
</method>
<method name="GetChat">
<arg name="session_id" type="s" direction="in"/>
<arg name="offset" type="x" direction="in"/>
<arg name="text" type="s" direction="out"/>
<arg name="new_offset" type="x" direction="out"/>
</method>
<method name="GetUsage">
<arg name="session_id" type="s" direction="in"/>
<arg name="usage" type="s" direction="out"/>
</method>
<method name="GetCost">
<arg name="session_id" type="s" direction="in"/>
<arg name="cost" type="s" direction="out"/>
</method>
<method name="GetConfig">
<arg name="session_id" type="s" direction="in"/>
<arg name="config" type="s" direction="out"/>
</method>
<method name="SetConfig">
<arg name="session_id" type="s" direction="in"/>
<arg name="key" type="s" direction="in"/>
<arg name="value" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="GetContext">
<arg name="session_id" type="s" direction="in"/>
<arg name="context" type="s" direction="out"/>
</method>
<method name="ListBackends">
<arg name="backends" type="as" direction="out"/>
</method>
<method name="ListModels">
<arg name="session_id" type="s" direction="in"/>
<arg name="models" type="as" direction="out"/>
</method>
<method name="ListAgents">
<arg name="agents" type="as" direction="out"/>
</method>
<method name="PeerAdd">
<arg name="session_id" type="s" direction="in"/>
<arg name="peer_id" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="PeerRemove">
<arg name="session_id" type="s" direction="in"/>
<arg name="peer_id" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="PeerList">
<arg name="session_id" type="s" direction="in"/>
<arg name="peers" type="as" direction="out"/>
</method>
<method name="PeerSubmit">
<arg name="session_id" type="s" direction="in"/>
<arg name="peer_id" type="s" direction="in"/>
<arg name="prompt" type="s" direction="in"/>
<arg name="success" type="b" direction="out"/>
</method>
<method name="Complete">
<arg name="cwd" type="s" direction="in"/>
<arg name="file" type="s" direction="in"/>
<arg name="prefix" type="s" direction="in"/>
<arg name="suffix" type="s" direction="in"/>
<arg name="context" type="s" direction="in"/>
<arg name="result" type="s" direction="out"/>
</method>
<signal name="SessionCreated">
<arg name="session_id" type="s"/>
</signal>
<signal name="SessionKilled">
<arg name="session_id" type="s"/>
</signal>
<signal name="SessionRenamed">
<arg name="old_id" type="s"/>
<arg name="new_id" type="s"/>
</signal>
<signal name="StateChanged">
<arg name="session_id" type="s"/>
<arg name="new_state" type="s"/>
</signal>
<signal name="ChatUpdated">
<arg name="session_id" type="s"/>
<arg name="offset" type="x"/>
<arg name="new_text" type="s"/>
</signal>
</interface>
` + introspect.IntrospectDataString + `
</node>`
func main() {
conn, err := dbus.ConnectSessionBus()
if err != nil {
fmt.Fprintf(os.Stderr, "cannot connect to session bus: %v\n", err)
os.Exit(1)
}
defer conn.Close()
reply, err := conn.RequestName(busName, dbus.NameFlagDoNotQueue)
if err != nil || reply != dbus.RequestNameReplyPrimaryOwner {
fmt.Fprintf(os.Stderr, "cannot acquire bus name %s: %v\n", busName, err)
os.Exit(1)
}
mgr := NewSessionManager(conn)
mgr.restoreAllSessions()
conn.Export(mgr, busPath, busIface)
conn.Export(introspect.Introspectable(introspectXML), busPath,
"org.freedesktop.DBus.Introspectable")
fmt.Println("ollie-dbus daemon started on D-Bus:", busName)
// Wait for signal
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
// Shutdown on signal or D-Bus disconnection (session bus gone).
select {
case <-sigCh:
case <-conn.Context().Done():
}
fmt.Println("ollie-dbus daemon shutting down")
mgr.Shutdown()
}