680 lines
16 KiB
Go
680 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
|
|
}
|
|
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="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()
|
|
}
|