diff --git a/.gitkeep b/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..9683649 --- /dev/null +++ b/Makefile @@ -0,0 +1,20 @@ +.PHONY: all build clean install uninstall + +BIN := ollied + +all: build + +build: + go build -o $(BIN) + +clean: + rm -f $(BIN) + +install: build + install -m755 $(BIN) /usr/bin/ + install -m644 org.ollie.SessionManager.desktop /etc/xdg/autostart/ + +uninstall: + pkill -u "$$(id -u)" -x $(BIN) 2>/dev/null; true + rm -f /usr/bin/$(BIN) + rm -f /etc/xdg/autostart/org.ollie.SessionManager.desktop diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..52d8e2b --- /dev/null +++ b/go.mod @@ -0,0 +1,15 @@ +module ollie-dbus + +go 1.25.6 + +require ( + github.com/godbus/dbus/v5 v5.1.0 + ollie v0.0.0-00010101000000-000000000000 +) + +require ( + github.com/simonfxr/pubsub v0.0.5 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) + +replace ollie => ../core diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..0847714 --- /dev/null +++ b/go.sum @@ -0,0 +1,16 @@ +github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/godbus/dbus/v5 v5.1.0 h1:4KLkAxT3aOY8Li4FRJe/KvhoNFFxo0m6fNuFUO8QJUk= +github.com/godbus/dbus/v5 v5.1.0/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/simonfxr/pubsub v0.0.5 h1:DJfvFoglqGvwJriIOC5NI5um34n2YX8KAU5+7jv768w= +github.com/simonfxr/pubsub v0.0.5/go.mod h1:bQ+B2NEEHZ08VY/0xVFGaGL1KMJ7cl3yLaXe8xvyC24= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.6.1 h1:hDPOHmpOpP40lSULcqw7IrRb/u7w6RpDC9399XyoNd0= +github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/main.go b/main.go new file mode 100644 index 0000000..4729a0f --- /dev/null +++ b/main.go @@ -0,0 +1,679 @@ +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 = ` + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + ` + introspect.IntrospectDataString + ` +` + +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() +} diff --git a/org.ollie.SessionManager.desktop b/org.ollie.SessionManager.desktop new file mode 100644 index 0000000..c365497 --- /dev/null +++ b/org.ollie.SessionManager.desktop @@ -0,0 +1,11 @@ +[Desktop Entry] +Type=Application +Name=Ollie Session Manager +GenericName=AI Agent D-Bus Daemon +Comment=D-Bus session manager for Ollie AI agent sessions +Exec=/usr/bin/ollied +Icon=ollie +Terminal=false +Categories=Utility; +X-DBUS-StartupType=unique +X-DBUS-ServiceName=org.ollie.SessionManager