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

685 lines
16 KiB
Go

package main
import (
"context"
"encoding/json"
"fmt"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"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
}
// 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})
// Create the agent core
core := agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: agentName,
AgentsDir: agentsDir,
SessionsDir: sessionsDir,
SessionID: sessID,
CWD: cwd,
Runtime: rt,
})
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
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
}
sess.logMu.Lock()
offset := int64(len(sess.log))
sess.log = append(sess.log, []byte(text)...)
sess.logMu.Unlock()
m.conn.Emit(busPath, busIface+".ChatUpdated", sess.id, offset, text)
})
// Watch state changes
go func() {
current := "idle"
for {
next, ok := core.WaitChange(ctx, agent.WatchState, current)
if !ok {
return
}
current = next
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.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.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
}
// --- 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
)
}
// --- Shutdown ---
func (m *SessionManager) Shutdown() {
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>
<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)
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)
<-sigCh
fmt.Println("ollie-dbus daemon shutting down")
mgr.Shutdown()
}