9p: rename session/ → mgr/, drop agent import alias
Rename the internal olliesrv/session package to olliesrv/mgr to resolve the naming conflict with ollie/session (the core session package). This allows importing ollie/session without an alias — the temporary `agent "ollie/session"` scaffolding is removed. All references now use the natural package name: session.X for core types, mgr.X for the 9P session manager.
This commit is contained in:
parent
4359f354af
commit
a6a7fae7ca
24
dbus.go
24
dbus.go
|
|
@ -13,9 +13,9 @@ import (
|
||||||
"github.com/godbus/dbus/v5"
|
"github.com/godbus/dbus/v5"
|
||||||
"github.com/godbus/dbus/v5/introspect"
|
"github.com/godbus/dbus/v5/introspect"
|
||||||
|
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
"olliesrv/session"
|
"olliesrv/mgr"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|
@ -24,10 +24,10 @@ const (
|
||||||
busIface = "org.ollie.SessionManager"
|
busIface = "org.ollie.SessionManager"
|
||||||
)
|
)
|
||||||
|
|
||||||
// DBusAdapter exposes a session.Manager over D-Bus.
|
// DBusAdapter exposes a mgr.Manager over D-Bus.
|
||||||
type DBusAdapter struct {
|
type DBusAdapter struct {
|
||||||
conn *dbus.Conn
|
conn *dbus.Conn
|
||||||
mgr *session.Manager
|
mgr *mgr.Manager
|
||||||
|
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
watchers map[string]context.CancelFunc // sessionID -> cancel for state watcher
|
watchers map[string]context.CancelFunc // sessionID -> cancel for state watcher
|
||||||
|
|
@ -35,7 +35,7 @@ type DBusAdapter struct {
|
||||||
|
|
||||||
// startDBus connects to the session bus, claims the well-known name, and
|
// startDBus connects to the session bus, claims the well-known name, and
|
||||||
// exports the adapter. Returns nil (no-op) if the bus is unavailable.
|
// exports the adapter. Returns nil (no-op) if the bus is unavailable.
|
||||||
func startDBus(mgr *session.Manager) *DBusAdapter {
|
func startDBus(mgr *mgr.Manager) *DBusAdapter {
|
||||||
conn, err := dbus.ConnectSessionBus()
|
conn, err := dbus.ConnectSessionBus()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Fprintf(os.Stderr, "dbus: session bus unavailable: %v\n", err)
|
fmt.Fprintf(os.Stderr, "dbus: session bus unavailable: %v\n", err)
|
||||||
|
|
@ -97,9 +97,9 @@ func (a *DBusAdapter) Close() {
|
||||||
a.conn.Close()
|
a.conn.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Lifecycle callbacks (called by session.Manager) ---
|
// --- Lifecycle callbacks (called by mgr.Manager) ---
|
||||||
|
|
||||||
func (a *DBusAdapter) OnSessionCreated(id string, sess *session.Session) {
|
func (a *DBusAdapter) OnSessionCreated(id string, sess *mgr.Session) {
|
||||||
if a == nil {
|
if a == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -129,7 +129,7 @@ func (a *DBusAdapter) OnSessionRenamed(oldID, newID string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// startWatcher launches goroutines for StateChanged and ChatUpdated signals.
|
// startWatcher launches goroutines for StateChanged and ChatUpdated signals.
|
||||||
func (a *DBusAdapter) startWatcher(id string, sess *session.Session) {
|
func (a *DBusAdapter) startWatcher(id string, sess *mgr.Session) {
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
a.mu.Lock()
|
a.mu.Lock()
|
||||||
a.watchers[id] = cancel
|
a.watchers[id] = cancel
|
||||||
|
|
@ -139,7 +139,7 @@ func (a *DBusAdapter) startWatcher(id string, sess *session.Session) {
|
||||||
go func() {
|
go func() {
|
||||||
current := sess.Core.State()
|
current := sess.Core.State()
|
||||||
for {
|
for {
|
||||||
next, ok := sess.Core.WaitChange(ctx, agent.WatchState, current)
|
next, ok := sess.Core.WaitChange(ctx, session.WatchState, current)
|
||||||
if !ok {
|
if !ok {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -155,7 +155,7 @@ func (a *DBusAdapter) startWatcher(id string, sess *session.Session) {
|
||||||
go func() {
|
go func() {
|
||||||
current := sess.Core.AgentName()
|
current := sess.Core.AgentName()
|
||||||
for {
|
for {
|
||||||
next, ok := sess.Core.WaitChange(ctx, agent.WatchAgent, current)
|
next, ok := sess.Core.WaitChange(ctx, session.WatchAgent, current)
|
||||||
if !ok {
|
if !ok {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -325,7 +325,7 @@ func (a *DBusAdapter) Interrupt(sessionID string) (bool, *dbus.Error) {
|
||||||
if sess == nil {
|
if sess == nil {
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
sess.Core.Interrupt(agent.ErrInterrupted)
|
sess.Core.Interrupt(session.ErrInterrupted)
|
||||||
return true, nil
|
return true, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -461,7 +461,7 @@ func (a *DBusAdapter) ListModels(sessionID string) ([]string, *dbus.Error) {
|
||||||
|
|
||||||
func (a *DBusAdapter) ListAgents() ([]string, *dbus.Error) {
|
func (a *DBusAdapter) ListAgents() ([]string, *dbus.Error) {
|
||||||
var result []string
|
var result []string
|
||||||
for _, dir := range agent.AgentsDirs() {
|
for _, dir := range session.AgentsDirs() {
|
||||||
entries, err := os.ReadDir(dir)
|
entries, err := os.ReadDir(dir)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
continue
|
continue
|
||||||
|
|
|
||||||
12
main.go
12
main.go
|
|
@ -14,7 +14,7 @@ import (
|
||||||
"syscall"
|
"syscall"
|
||||||
|
|
||||||
"9fans.net/go/plan9/client"
|
"9fans.net/go/plan9/client"
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
"ollie/elevate"
|
"ollie/elevate"
|
||||||
"ollie/env"
|
"ollie/env"
|
||||||
|
|
@ -25,7 +25,7 @@ import (
|
||||||
fs "olliesrv/fs"
|
fs "olliesrv/fs"
|
||||||
"olliesrv/mount"
|
"olliesrv/mount"
|
||||||
"olliesrv/server"
|
"olliesrv/server"
|
||||||
"olliesrv/session"
|
"olliesrv/mgr"
|
||||||
)
|
)
|
||||||
|
|
||||||
const serviceName = "ollie"
|
const serviceName = "ollie"
|
||||||
|
|
@ -157,7 +157,7 @@ func runServer(sockPath string) {
|
||||||
|
|
||||||
sink := olog.NewSink(os.Stdout, os.Stderr, olog.ParseLevel(os.Getenv("OLLIE_LOG"), olog.LevelWarn))
|
sink := olog.NewSink(os.Stdout, os.Stderr, olog.ParseLevel(os.Getenv("OLLIE_LOG"), olog.LevelWarn))
|
||||||
|
|
||||||
agentsDirs := agent.AgentsDirs()
|
agentsDirs := session.AgentsDirs()
|
||||||
sessionsDir := paths.DataDir() + "/sessions"
|
sessionsDir := paths.DataDir() + "/sessions"
|
||||||
|
|
||||||
// Create the tool registry
|
// Create the tool registry
|
||||||
|
|
@ -178,7 +178,7 @@ func runServer(sockPath string) {
|
||||||
// Elevate broker (initialized after manager; closures capture the pointer).
|
// Elevate broker (initialized after manager; closures capture the pointer).
|
||||||
var elevateBroker *elevate.Broker
|
var elevateBroker *elevate.Broker
|
||||||
|
|
||||||
mgr := session.NewManager(session.ManagerConfig{
|
mgr := mgr.NewManager(mgr.ManagerConfig{
|
||||||
ToolRegistry: toolRegistry,
|
ToolRegistry: toolRegistry,
|
||||||
SkillsRegistry: skillsRegistry,
|
SkillsRegistry: skillsRegistry,
|
||||||
AgentsDir: agentsDirs[0],
|
AgentsDir: agentsDirs[0],
|
||||||
|
|
@ -196,7 +196,7 @@ func runServer(sockPath string) {
|
||||||
elevateBroker.ResetTurn(sessionID)
|
elevateBroker.ResetTurn(sessionID)
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
OnSessionCreated: func(id string, sess *session.Session) {
|
OnSessionCreated: func(id string, sess *mgr.Session) {
|
||||||
if dbusAdapter != nil {
|
if dbusAdapter != nil {
|
||||||
dbusAdapter.OnSessionCreated(id, sess)
|
dbusAdapter.OnSessionCreated(id, sess)
|
||||||
}
|
}
|
||||||
|
|
@ -366,7 +366,7 @@ func NewRootStore(reg *tools.Registry) *fs.Tree {
|
||||||
},
|
},
|
||||||
"agents": func() ([]byte, error) {
|
"agents": func() ([]byte, error) {
|
||||||
var sb strings.Builder
|
var sb strings.Builder
|
||||||
for _, dir := range agent.AgentsDirs() {
|
for _, dir := range session.AgentsDirs() {
|
||||||
entries, err := os.ReadDir(dir)
|
entries, err := os.ReadDir(dir)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
continue
|
continue
|
||||||
|
|
|
||||||
134
main_test.go
134
main_test.go
|
|
@ -11,21 +11,21 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
olog "ollie/log"
|
olog "ollie/log"
|
||||||
"olliesrv/fs"
|
"olliesrv/fs"
|
||||||
"olliesrv/session"
|
"olliesrv/mgr"
|
||||||
)
|
)
|
||||||
|
|
||||||
// --- test helpers ---
|
// --- test helpers ---
|
||||||
|
|
||||||
func testSink() *olog.Sink { return olog.NewSink(io.Discard, io.Discard, olog.LevelError) }
|
func testSink() *olog.Sink { return olog.NewSink(io.Discard, io.Discard, olog.LevelError) }
|
||||||
|
|
||||||
// newNoopCore creates a real *agent.Session with a noop backend.
|
// newNoopCore creates a real *session.Session with a noop backend.
|
||||||
func newNoopCore(id string) *agent.Session {
|
func newNoopCore(id string) *session.Session {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: "default",
|
AgentName: "default",
|
||||||
CWD: "/tmp",
|
CWD: "/tmp",
|
||||||
|
|
@ -36,8 +36,8 @@ func newNoopCore(id string) *agent.Session {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// newBlockingCore creates a real *agent.Session whose backend blocks until ctx is cancelled.
|
// newBlockingCore creates a real *session.Session whose backend blocks until ctx is cancelled.
|
||||||
func newBlockingCore(id string) *agent.Session {
|
func newBlockingCore(id string) *session.Session {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
be.ChatStreamFunc = func(ctx context.Context, _ []backend.Message, _ []backend.Tool, _ backend.GenerationParams) (<-chan backend.StreamEvent, error) {
|
be.ChatStreamFunc = func(ctx context.Context, _ []backend.Message, _ []backend.Tool, _ backend.GenerationParams) (<-chan backend.StreamEvent, error) {
|
||||||
ch := make(chan backend.StreamEvent, 1)
|
ch := make(chan backend.StreamEvent, 1)
|
||||||
|
|
@ -48,7 +48,7 @@ func newBlockingCore(id string) *agent.Session {
|
||||||
}()
|
}()
|
||||||
return ch, nil
|
return ch, nil
|
||||||
}
|
}
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: "default",
|
AgentName: "default",
|
||||||
CWD: "/tmp",
|
CWD: "/tmp",
|
||||||
|
|
@ -59,8 +59,8 @@ func newBlockingCore(id string) *agent.Session {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// newContentCore creates a real *agent.Session whose backend emits specific content.
|
// newContentCore creates a real *session.Session whose backend emits specific content.
|
||||||
func newContentCore(id, content string) *agent.Session {
|
func newContentCore(id, content string) *session.Session {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
be.ChatStreamFunc = func(ctx context.Context, _ []backend.Message, _ []backend.Tool, _ backend.GenerationParams) (<-chan backend.StreamEvent, error) {
|
be.ChatStreamFunc = func(ctx context.Context, _ []backend.Message, _ []backend.Tool, _ backend.GenerationParams) (<-chan backend.StreamEvent, error) {
|
||||||
ch := make(chan backend.StreamEvent, 2)
|
ch := make(chan backend.StreamEvent, 2)
|
||||||
|
|
@ -69,7 +69,7 @@ func newContentCore(id, content string) *agent.Session {
|
||||||
close(ch)
|
close(ch)
|
||||||
return ch, nil
|
return ch, nil
|
||||||
}
|
}
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: "default",
|
AgentName: "default",
|
||||||
CWD: "/tmp",
|
CWD: "/tmp",
|
||||||
|
|
@ -80,22 +80,22 @@ func newContentCore(id, content string) *agent.Session {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
func testSession(id string) *session.Session {
|
func testSession(id string) *mgr.Session {
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
core := newNoopCore(id)
|
core := newNoopCore(id)
|
||||||
return session.NewSession(id, core, ctx, cancel)
|
return mgr.NewSession(id, core, ctx, cancel)
|
||||||
}
|
}
|
||||||
|
|
||||||
func testBlockingSession(id string) *session.Session {
|
func testBlockingSession(id string) *mgr.Session {
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
core := newBlockingCore(id)
|
core := newBlockingCore(id)
|
||||||
return session.NewSession(id, core, ctx, cancel)
|
return mgr.NewSession(id, core, ctx, cancel)
|
||||||
}
|
}
|
||||||
|
|
||||||
func newTestSessionManager(t *testing.T) *session.Manager {
|
func newTestSessionManager(t *testing.T) *mgr.Manager {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
return session.NewManager(session.ManagerConfig{
|
return mgr.NewManager(mgr.ManagerConfig{
|
||||||
Log: sink.NewLogger("test"),
|
Log: sink.NewLogger("test"),
|
||||||
Sink: sink,
|
Sink: sink,
|
||||||
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
||||||
|
|
@ -103,17 +103,17 @@ func newTestSessionManager(t *testing.T) *session.Manager {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
func newTestSessionManagerWithCore(t *testing.T) *session.Manager {
|
func newTestSessionManagerWithCore(t *testing.T) *mgr.Manager {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
return session.NewManager(session.ManagerConfig{
|
return mgr.NewManager(mgr.ManagerConfig{
|
||||||
Log: sink.NewLogger("test"),
|
Log: sink.NewLogger("test"),
|
||||||
Sink: sink,
|
Sink: sink,
|
||||||
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
||||||
MkdirAll: func(string, os.FileMode) error { return nil },
|
MkdirAll: func(string, os.FileMode) error { return nil },
|
||||||
NewCore: func(sessionID, agentName, cwd string) (*agent.Session, error) {
|
NewCore: func(sessionID, agentName, cwd string) (*session.Session, error) {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: agentName,
|
AgentName: agentName,
|
||||||
CWD: cwd,
|
CWD: cwd,
|
||||||
|
|
@ -157,11 +157,11 @@ var _ = os.Remove // suppress unused import
|
||||||
var _ = filepath.Join // suppress unused import
|
var _ = filepath.Join // suppress unused import
|
||||||
var _ sync.Mutex // suppress unused import
|
var _ sync.Mutex // suppress unused import
|
||||||
|
|
||||||
func newTestSessionFileStore(t *testing.T, sess *session.Session) *fs.Tree {
|
func newTestSessionFileStore(t *testing.T, sess *mgr.Session) *fs.Tree {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
// Use agent tree for agent-level file tests
|
// Use agent tree for agent-level file tests
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func(data []byte) error { return nil },
|
func(data []byte) error { return nil },
|
||||||
func() {},
|
func() {},
|
||||||
nil,
|
nil,
|
||||||
|
|
@ -170,14 +170,14 @@ func newTestSessionFileStore(t *testing.T, sess *session.Session) *fs.Tree {
|
||||||
return sf
|
return sf
|
||||||
}
|
}
|
||||||
|
|
||||||
func newTestSessionFileStoreWith(t *testing.T, sess *session.Session, kill func(), rename func(string) error, save func([]byte) error) *fs.Tree {
|
func newTestSessionFileStoreWith(t *testing.T, sess *mgr.Session, kill func(), rename func(string) error, save func([]byte) error) *fs.Tree {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
// Session tree for kill/rename/save operations
|
// Session tree for kill/rename/save operations
|
||||||
return session.NewSessionTree(sess, sink.NewLogger("test"), kill, rename, save, func() {}, nil, nil)
|
return mgr.NewSessionTree(sess, sink.NewLogger("test"), kill, rename, save, func() {}, nil, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ===== session.Session =====
|
// ===== mgr.Session =====
|
||||||
|
|
||||||
func TestSessionAppendLog(t *testing.T) {
|
func TestSessionAppendLog(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
|
|
@ -197,15 +197,15 @@ func TestSessionAppendLog(t *testing.T) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestSessionManagerFileMode(t *testing.T) {
|
func TestSessionManagerFileMode(t *testing.T) {
|
||||||
if m, ok := session.FileMode("new"); !ok || m != 0666 {
|
if m, ok := mgr.FileMode("new"); !ok || m != 0666 {
|
||||||
t.Errorf("session.FileMode(new) = %o, %v", m, ok)
|
t.Errorf("mgr.FileMode(new) = %o, %v", m, ok)
|
||||||
}
|
}
|
||||||
if _, ok := session.FileMode("bogus"); ok {
|
if _, ok := mgr.FileMode("bogus"); ok {
|
||||||
t.Error("session.FileMode(bogus) should be false")
|
t.Error("mgr.FileMode(bogus) should be false")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ===== session.Manager =====
|
// ===== mgr.Manager =====
|
||||||
|
|
||||||
func TestSessionManagerReadableContract(t *testing.T) {
|
func TestSessionManagerReadableContract(t *testing.T) {
|
||||||
}
|
}
|
||||||
|
|
@ -304,13 +304,13 @@ func TestSessionManagerDeleteAndKill(t *testing.T) {
|
||||||
func TestSessionManagerSession(t *testing.T) {
|
func TestSessionManagerSession(t *testing.T) {
|
||||||
s := newTestSessionManager(t)
|
s := newTestSessionManager(t)
|
||||||
if s.Session("nope") != nil {
|
if s.Session("nope") != nil {
|
||||||
t.Error("session.Session(nonexistent) should be nil")
|
t.Error("mgr.Session(nonexistent) should be nil")
|
||||||
}
|
}
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
s.AddSession(sess)
|
s.AddSession(sess)
|
||||||
if s.Session("s1") == nil {
|
if s.Session("s1") == nil {
|
||||||
t.Error("session.Session(s1) should not be nil")
|
t.Error("mgr.Session(s1) should not be nil")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -340,7 +340,7 @@ func TestSessionManagerShutdown(t *testing.T) {
|
||||||
|
|
||||||
func TestSessionManagerRename(t *testing.T) {
|
func TestSessionManagerRename(t *testing.T) {
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
s := session.NewManager(session.ManagerConfig{
|
s := mgr.NewManager(mgr.ManagerConfig{
|
||||||
Log: sink.NewLogger("test"),
|
Log: sink.NewLogger("test"),
|
||||||
Sink: sink,
|
Sink: sink,
|
||||||
ReadFile: func(string) ([]byte, error) { return nil, nil },
|
ReadFile: func(string) ([]byte, error) { return nil, nil },
|
||||||
|
|
@ -390,7 +390,7 @@ func TestSessionFileStoreReadableContract(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
_ = session.NewAgentTree(sess, sink.NewLogger("test"),
|
_ = mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -398,15 +398,15 @@ func TestSessionFileStoreList(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
entries, err := sf.List()
|
entries, err := sf.List()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("List: %v", err)
|
t.Fatalf("List: %v", err)
|
||||||
}
|
}
|
||||||
if len(entries) != len(session.AgentFileList) {
|
if len(entries) != len(mgr.AgentFileList) {
|
||||||
t.Errorf("List() returned %d entries; want %d", len(entries), len(session.AgentFileList))
|
t.Errorf("List() returned %d entries; want %d", len(entries), len(mgr.AgentFileList))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -415,7 +415,7 @@ func TestSessionFileStoreStatChat(t *testing.T) {
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sess.AppendLog([]byte("hello"))
|
sess.AppendLog([]byte("hello"))
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
fi, err := sf.Stat("chat")
|
fi, err := sf.Stat("chat")
|
||||||
|
|
@ -432,7 +432,7 @@ func TestSessionFileStoreGetChat(t *testing.T) {
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sess.AppendLog([]byte("hello"))
|
sess.AppendLog([]byte("hello"))
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
data := testStoreRead(t, sf, "chat")
|
data := testStoreRead(t, sf, "chat")
|
||||||
|
|
@ -445,7 +445,7 @@ func TestSessionFileStoreGetContent(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
for _, name := range []string{"cfg", "offset", "usage", "ctxsz", "models", "systemprompt"} {
|
for _, name := range []string{"cfg", "offset", "usage", "ctxsz", "models", "systemprompt"} {
|
||||||
|
|
@ -459,7 +459,7 @@ func TestSessionFileStorePutCwd(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
// Create the directory first so SetCWD validates it
|
// Create the directory first so SetCWD validates it
|
||||||
|
|
@ -475,7 +475,7 @@ func TestSessionFileStorePutEmpty(t *testing.T) {
|
||||||
sess := testSession("s1")
|
sess := testSession("s1")
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func([]byte) error { return nil }, func() {}, nil, nil)
|
func([]byte) error { return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
// Empty write is a no-op
|
// Empty write is a no-op
|
||||||
|
|
@ -582,7 +582,7 @@ func TestSessionFileStoreWriteChat(t *testing.T) {
|
||||||
defer sess.Cancel()
|
defer sess.Cancel()
|
||||||
var saved []byte
|
var saved []byte
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
sf := session.NewAgentTree(sess, sink.NewLogger("test"),
|
sf := mgr.NewAgentTree(sess, sink.NewLogger("test"),
|
||||||
func(data []byte) error { saved = data; return nil }, func() {}, nil, nil)
|
func(data []byte) error { saved = data; return nil }, func() {}, nil, nil)
|
||||||
|
|
||||||
testStoreWrite(t, sf, "chat", []byte("transcript data"))
|
testStoreWrite(t, sf, "chat", []byte("transcript data"))
|
||||||
|
|
@ -824,7 +824,7 @@ func TestSessionFileStoreMakePublish(t *testing.T) {
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
core := newContentCore("s1", "hello world")
|
core := newContentCore("s1", "hello world")
|
||||||
sess := session.NewSession("s1", core, ctx, cancel)
|
sess := mgr.NewSession("s1", core, ctx, cancel)
|
||||||
sf := newTestSessionFileStore(t, sess)
|
sf := newTestSessionFileStore(t, sess)
|
||||||
|
|
||||||
// Submit triggers the content backend
|
// Submit triggers the content backend
|
||||||
|
|
@ -850,7 +850,7 @@ func TestSessionFileStoreMakePublishMultipleEvents(t *testing.T) {
|
||||||
close(ch)
|
close(ch)
|
||||||
return ch, nil
|
return ch, nil
|
||||||
}
|
}
|
||||||
core := agent.New(agent.Config{
|
core := session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: "default",
|
AgentName: "default",
|
||||||
CWD: "/tmp",
|
CWD: "/tmp",
|
||||||
|
|
@ -858,7 +858,7 @@ func TestSessionFileStoreMakePublishMultipleEvents(t *testing.T) {
|
||||||
})
|
})
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
sess := session.NewSession("s1", core, ctx, cancel)
|
sess := mgr.NewSession("s1", core, ctx, cancel)
|
||||||
sf := newTestSessionFileStore(t, sess)
|
sf := newTestSessionFileStore(t, sess)
|
||||||
|
|
||||||
testStoreWrite(t, sf, "prompt", []byte("multi"))
|
testStoreWrite(t, sf, "prompt", []byte("multi"))
|
||||||
|
|
@ -931,7 +931,7 @@ func TestLoadAgentConfig(t *testing.T) {
|
||||||
data := []byte(`{"prompt":"test prompt","maxTokens":1024}`)
|
data := []byte(`{"prompt":"test prompt","maxTokens":1024}`)
|
||||||
os.WriteFile(filepath.Join(dir, "test.json"), data, 0644)
|
os.WriteFile(filepath.Join(dir, "test.json"), data, 0644)
|
||||||
|
|
||||||
cfg := session.LoadAgentConfig(dir, "test", nil)
|
cfg := mgr.LoadAgentConfig(dir, "test", nil)
|
||||||
if cfg == nil {
|
if cfg == nil {
|
||||||
t.Fatal("LoadAgentConfig returned nil")
|
t.Fatal("LoadAgentConfig returned nil")
|
||||||
}
|
}
|
||||||
|
|
@ -942,17 +942,17 @@ func TestLoadAgentConfig(t *testing.T) {
|
||||||
|
|
||||||
func TestFormatEvent(t *testing.T) {
|
func TestFormatEvent(t *testing.T) {
|
||||||
tests := []struct {
|
tests := []struct {
|
||||||
ev agent.Event
|
ev session.Event
|
||||||
want string
|
want string
|
||||||
}{
|
}{
|
||||||
{agent.Event{Role: "user", Content: "hello"}, "[user]\nhello\n"},
|
{session.Event{Role: "user", Content: "hello"}, "[user]\nhello\n"},
|
||||||
{agent.Event{Role: "assistant", Content: "hi"}, "hi"},
|
{session.Event{Role: "assistant", Content: "hi"}, "hi"},
|
||||||
{agent.Event{Role: "call", Name: "fn", Content: "args"}, "[call:fn]\nargs\n"},
|
{session.Event{Role: "call", Name: "fn", Content: "args"}, "[call:fn]\nargs\n"},
|
||||||
{agent.Event{Role: "tool", Name: "fn", Content: "result"}, "[tool:fn]\nresult\n"},
|
{session.Event{Role: "tool", Name: "fn", Content: "result"}, "[tool:fn]\nresult\n"},
|
||||||
{agent.Event{Role: "info", Content: "msg\n"}, "[info]\nmsg\n"},
|
{session.Event{Role: "info", Content: "msg\n"}, "[info]\nmsg\n"},
|
||||||
}
|
}
|
||||||
for _, tc := range tests {
|
for _, tc := range tests {
|
||||||
got := string(session.FormatEvent(tc.ev))
|
got := string(mgr.FormatEvent(tc.ev))
|
||||||
if got != tc.want {
|
if got != tc.want {
|
||||||
t.Errorf("FormatEvent(%v) = %q; want %q", tc.ev.Role, got, tc.want)
|
t.Errorf("FormatEvent(%v) = %q; want %q", tc.ev.Role, got, tc.want)
|
||||||
}
|
}
|
||||||
|
|
@ -997,7 +997,7 @@ func TestFormatParamsRoundTrip(t *testing.T) {
|
||||||
Temperature: &temp,
|
Temperature: &temp,
|
||||||
TopP: &topP,
|
TopP: &topP,
|
||||||
}
|
}
|
||||||
text := session.FormatParams(p)
|
text := mgr.FormatParams(p)
|
||||||
if !strings.Contains(text, "maxTokens=4096") {
|
if !strings.Contains(text, "maxTokens=4096") {
|
||||||
t.Errorf("FormatParams missing maxTokens; got: %s", text)
|
t.Errorf("FormatParams missing maxTokens; got: %s", text)
|
||||||
}
|
}
|
||||||
|
|
@ -1005,7 +1005,7 @@ func TestFormatParamsRoundTrip(t *testing.T) {
|
||||||
t.Errorf("FormatParams missing temperature; got: %s", text)
|
t.Errorf("FormatParams missing temperature; got: %s", text)
|
||||||
}
|
}
|
||||||
|
|
||||||
parsed, err := session.ParseParams(text, backend.GenerationParams{})
|
parsed, err := mgr.ParseParams(text, backend.GenerationParams{})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("ParseParams: %v", err)
|
t.Fatalf("ParseParams: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -1019,7 +1019,7 @@ func TestFormatParamsRoundTrip(t *testing.T) {
|
||||||
|
|
||||||
func TestFormatParamsNilOptionals(t *testing.T) {
|
func TestFormatParamsNilOptionals(t *testing.T) {
|
||||||
p := backend.GenerationParams{MaxTokens: 100}
|
p := backend.GenerationParams{MaxTokens: 100}
|
||||||
text := session.FormatParams(p)
|
text := mgr.FormatParams(p)
|
||||||
// Real FormatParams lists all fields; nil optionals appear as empty values
|
// Real FormatParams lists all fields; nil optionals appear as empty values
|
||||||
if !strings.Contains(text, "maxTokens=100") {
|
if !strings.Contains(text, "maxTokens=100") {
|
||||||
t.Errorf("FormatParams missing maxTokens=100; got: %s", text)
|
t.Errorf("FormatParams missing maxTokens=100; got: %s", text)
|
||||||
|
|
@ -1033,7 +1033,7 @@ func TestFormatParamsNilOptionals(t *testing.T) {
|
||||||
func TestParseParamsClearWithEmpty(t *testing.T) {
|
func TestParseParamsClearWithEmpty(t *testing.T) {
|
||||||
temp := 0.7
|
temp := 0.7
|
||||||
base := backend.GenerationParams{Temperature: &temp, MaxTokens: 100}
|
base := backend.GenerationParams{Temperature: &temp, MaxTokens: 100}
|
||||||
parsed, err := session.ParseParams("temperature=\nmaxTokens=200", base)
|
parsed, err := mgr.ParseParams("temperature=\nmaxTokens=200", base)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("ParseParams: %v", err)
|
t.Fatalf("ParseParams: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -1046,11 +1046,11 @@ func TestParseParamsClearWithEmpty(t *testing.T) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestParseParamsErrors(t *testing.T) {
|
func TestParseParamsErrors(t *testing.T) {
|
||||||
_, err := session.ParseParams("maxTokens=notanumber", backend.GenerationParams{})
|
_, err := mgr.ParseParams("maxTokens=notanumber", backend.GenerationParams{})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Error("non-numeric maxTokens should error")
|
t.Error("non-numeric maxTokens should error")
|
||||||
}
|
}
|
||||||
_, err = session.ParseParams("temperature=notafloat", backend.GenerationParams{})
|
_, err = mgr.ParseParams("temperature=notafloat", backend.GenerationParams{})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Error("non-numeric temperature should error")
|
t.Error("non-numeric temperature should error")
|
||||||
}
|
}
|
||||||
|
|
@ -1137,14 +1137,14 @@ func TestSessionManagerCreateSessionUnknownKey(t *testing.T) {
|
||||||
func TestSessionManagerCreateSessionEnvExpansion(t *testing.T) {
|
func TestSessionManagerCreateSessionEnvExpansion(t *testing.T) {
|
||||||
t.Setenv("TEST_CWD", "/expanded/path")
|
t.Setenv("TEST_CWD", "/expanded/path")
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
s := session.NewManager(session.ManagerConfig{
|
s := mgr.NewManager(mgr.ManagerConfig{
|
||||||
Log: sink.NewLogger("test"),
|
Log: sink.NewLogger("test"),
|
||||||
Sink: sink,
|
Sink: sink,
|
||||||
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
||||||
MkdirAll: func(string, os.FileMode) error { return nil },
|
MkdirAll: func(string, os.FileMode) error { return nil },
|
||||||
NewCore: func(sessionID, agentName, cwd string) (*agent.Session, error) {
|
NewCore: func(sessionID, agentName, cwd string) (*session.Session, error) {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: agentName,
|
AgentName: agentName,
|
||||||
CWD: cwd,
|
CWD: cwd,
|
||||||
|
|
@ -1160,14 +1160,14 @@ func TestSessionManagerCreateSessionEnvExpansion(t *testing.T) {
|
||||||
|
|
||||||
func TestSessionManagerCreateSessionTildeExpansion(t *testing.T) {
|
func TestSessionManagerCreateSessionTildeExpansion(t *testing.T) {
|
||||||
sink := testSink()
|
sink := testSink()
|
||||||
s := session.NewManager(session.ManagerConfig{
|
s := mgr.NewManager(mgr.ManagerConfig{
|
||||||
Log: sink.NewLogger("test"),
|
Log: sink.NewLogger("test"),
|
||||||
Sink: sink,
|
Sink: sink,
|
||||||
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
ReadFile: func(string) ([]byte, error) { return []byte("#!/bin/sh\n"), nil },
|
||||||
MkdirAll: func(string, os.FileMode) error { return nil },
|
MkdirAll: func(string, os.FileMode) error { return nil },
|
||||||
NewCore: func(sessionID, agentName, cwd string) (*agent.Session, error) {
|
NewCore: func(sessionID, agentName, cwd string) (*session.Session, error) {
|
||||||
be := backend.NewNoop("stub", "m")
|
be := backend.NewNoop("stub", "m")
|
||||||
return agent.New(agent.Config{
|
return session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: agentName,
|
AgentName: agentName,
|
||||||
CWD: cwd,
|
CWD: cwd,
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
|
|
@ -64,7 +64,7 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
||||||
|
|
||||||
sessID := name
|
sessID := name
|
||||||
if sessID == "" {
|
if sessID == "" {
|
||||||
sessID = agent.NewSessionID()
|
sessID = session.NewSessionID()
|
||||||
}
|
}
|
||||||
|
|
||||||
s.mu.RLock()
|
s.mu.RLock()
|
||||||
|
|
@ -74,7 +74,7 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
||||||
return "", fmt.Errorf("session already exists: %s", sessID)
|
return "", fmt.Errorf("session already exists: %s", sessID)
|
||||||
}
|
}
|
||||||
|
|
||||||
var core *agent.Session
|
var core *session.Session
|
||||||
var sessPtr *Session
|
var sessPtr *Session
|
||||||
uname := s.nextUname()
|
uname := s.nextUname()
|
||||||
if s.cfg.NewCore != nil {
|
if s.cfg.NewCore != nil {
|
||||||
|
|
@ -161,7 +161,7 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
||||||
if len(remoteEnv) > 0 {
|
if len(remoteEnv) > 0 {
|
||||||
promptEnv = remoteEnv
|
promptEnv = remoteEnv
|
||||||
} else {
|
} else {
|
||||||
promptEnv = agent.PromptEnv(cwd)
|
promptEnv = session.PromptEnv(cwd)
|
||||||
}
|
}
|
||||||
env := []string{"OLLIE_SESSION_ID=" + sessID, "OLLIE_UNAME=" + uname}
|
env := []string{"OLLIE_SESSION_ID=" + sessID, "OLLIE_UNAME=" + uname}
|
||||||
env = append(env, promptEnv...)
|
env = append(env, promptEnv...)
|
||||||
|
|
@ -197,10 +197,10 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
||||||
envBlock := prompts.Environment(cwd, platform, isGitRepo, "")
|
envBlock := prompts.Environment(cwd, platform, isGitRepo, "")
|
||||||
|
|
||||||
disp := newDisp()
|
disp := newDisp()
|
||||||
rt := agent.BuildRuntime(cfg, disp, cwd, env, sysPrompt, opModel, envBlock)
|
rt := session.BuildRuntime(cfg, disp, cwd, env, sysPrompt, opModel, envBlock)
|
||||||
|
|
||||||
// sessPtr is set after NewSession; the ReadPlanStep closure captures it.
|
// sessPtr is set after NewSession; the ReadPlanStep closure captures it.
|
||||||
core = agent.New(agent.Config{
|
core = session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: agentName,
|
AgentName: agentName,
|
||||||
AgentsDir: s.cfg.AgentsDir,
|
AgentsDir: s.cfg.AgentsDir,
|
||||||
|
|
@ -222,7 +222,7 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
||||||
data := make([]byte, len(sessPtr.plan))
|
data := make([]byte, len(sessPtr.plan))
|
||||||
copy(data, sessPtr.plan)
|
copy(data, sessPtr.plan)
|
||||||
sessPtr.mu.RUnlock()
|
sessPtr.mu.RUnlock()
|
||||||
return agent.NextUncheckedStep(data)
|
return session.NextUncheckedStep(data)
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -1,30 +1,30 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
)
|
)
|
||||||
|
|
||||||
func LoadAgentConfig(agentsDir, name string, open func(string) (*os.File, error)) *agent.AgentConfig {
|
func LoadAgentConfig(agentsDir, name string, open func(string) (*os.File, error)) *session.AgentConfig {
|
||||||
if open == nil {
|
if open == nil {
|
||||||
open = os.Open
|
open = os.Open
|
||||||
}
|
}
|
||||||
path := agent.AgentConfigPath(agentsDir, name)
|
path := session.AgentConfigPath(agentsDir, name)
|
||||||
f, err := open(path)
|
f, err := open(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
defer f.Close()
|
defer f.Close()
|
||||||
cfg, _ := agent.Load(f)
|
cfg, _ := session.Load(f)
|
||||||
return cfg
|
return cfg
|
||||||
}
|
}
|
||||||
|
|
||||||
// FormatEvent converts an agent Event to bytes for appending to a chat log.
|
// FormatEvent converts an agent Event to bytes for appending to a chat log.
|
||||||
// Streaming roles (assistant, reasoning) return only the content chunk;
|
// Streaming roles (assistant, reasoning) return only the content chunk;
|
||||||
// the caller (startEventLog) is responsible for writing the [role] header.
|
// the caller (startEventLog) is responsible for writing the [role] header.
|
||||||
func FormatEvent(ev agent.Event) []byte {
|
func FormatEvent(ev session.Event) []byte {
|
||||||
switch ev.Role {
|
switch ev.Role {
|
||||||
case "user":
|
case "user":
|
||||||
return []byte("[user]\n" + ev.Content + "\n")
|
return []byte("[user]\n" + ev.Content + "\n")
|
||||||
|
|
@ -1,11 +1,11 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
olog "ollie/log"
|
olog "ollie/log"
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"ollie/paths"
|
"ollie/paths"
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/skills"
|
"ollie/skills"
|
||||||
"ollie/tools"
|
"ollie/tools"
|
||||||
"olliesrv/fs"
|
"olliesrv/fs"
|
||||||
|
|
@ -36,7 +36,7 @@ type ManagerConfig struct {
|
||||||
MkdirAll func(string, os.FileMode) error
|
MkdirAll func(string, os.FileMode) error
|
||||||
// NewCore, if non-nil, replaces the default backend.New + agent.New
|
// NewCore, if non-nil, replaces the default backend.New + agent.New
|
||||||
// path. It receives the session ID, agent name, and cwd, and returns a Session.
|
// path. It receives the session ID, agent name, and cwd, and returns a Session.
|
||||||
NewCore func(sessionID, agentName, cwd string) (*agent.Session, error)
|
NewCore func(sessionID, agentName, cwd string) (*session.Session, error)
|
||||||
// Strict rejects inline code steps; only tool steps are allowed.
|
// Strict rejects inline code steps; only tool steps are allowed.
|
||||||
Strict bool
|
Strict bool
|
||||||
// Yolo skips the landrun sandbox.
|
// Yolo skips the landrun sandbox.
|
||||||
|
|
@ -588,7 +588,7 @@ func (s *Manager) InterruptAll() {
|
||||||
s.mu.RLock()
|
s.mu.RLock()
|
||||||
defer s.mu.RUnlock()
|
defer s.mu.RUnlock()
|
||||||
for _, sess := range s.sessions {
|
for _, sess := range s.sessions {
|
||||||
sess.Core.Interrupt(agent.ErrInterrupted)
|
sess.Core.Interrupt(session.ErrInterrupted)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
|
|
@ -63,7 +63,7 @@ func (s *Manager) restoreAllSessions() {
|
||||||
|
|
||||||
// Load all persisted session JSONs (fast, sequential disk reads)
|
// Load all persisted session JSONs (fast, sequential disk reads)
|
||||||
type loadedSession struct {
|
type loadedSession struct {
|
||||||
ps *agent.PersistedAgent
|
ps *session.PersistedAgent
|
||||||
name string
|
name string
|
||||||
}
|
}
|
||||||
var loaded []loadedSession
|
var loaded []loadedSession
|
||||||
|
|
@ -72,7 +72,7 @@ func (s *Manager) restoreAllSessions() {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
path := filepath.Join(dir, e.Name())
|
path := filepath.Join(dir, e.Name())
|
||||||
ps, err := agent.LoadPersistedAgent(path)
|
ps, err := session.LoadPersistedAgent(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.cfg.Log.Error("restore session %s: %v", e.Name(), err)
|
s.cfg.Log.Error("restore session %s: %v", e.Name(), err)
|
||||||
continue
|
continue
|
||||||
|
|
@ -98,7 +98,7 @@ func (s *Manager) restoreAllSessions() {
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Manager) restoreSession(ps *agent.PersistedAgent) error {
|
func (s *Manager) restoreSession(ps *session.PersistedAgent) error {
|
||||||
cwd := ps.CWD
|
cwd := ps.CWD
|
||||||
if cwd == "" {
|
if cwd == "" {
|
||||||
cwd, _ = os.Getwd()
|
cwd, _ = os.Getwd()
|
||||||
|
|
@ -170,7 +170,7 @@ func (s *Manager) restoreSession(ps *agent.PersistedAgent) error {
|
||||||
newDisp = tools.NewDispatcherFunc(map[string]func() tools.Server{
|
newDisp = tools.NewDispatcherFunc(map[string]func() tools.Server{
|
||||||
"execute": execute.Decl(cwd, execOpts...),
|
"execute": execute.Decl(cwd, execOpts...),
|
||||||
})
|
})
|
||||||
promptEnv = agent.PromptEnv(cwd)
|
promptEnv = session.PromptEnv(cwd)
|
||||||
}
|
}
|
||||||
|
|
||||||
env := []string{"OLLIE_SESSION_ID=" + sessID, "OLLIE_UNAME=" + uname}
|
env := []string{"OLLIE_SESSION_ID=" + sessID, "OLLIE_UNAME=" + uname}
|
||||||
|
|
@ -204,12 +204,12 @@ func (s *Manager) restoreSession(ps *agent.PersistedAgent) error {
|
||||||
envBlock := prompts.Environment(cwd, platform, isGitRepo, "")
|
envBlock := prompts.Environment(cwd, platform, isGitRepo, "")
|
||||||
|
|
||||||
disp := newDisp()
|
disp := newDisp()
|
||||||
rt := agent.BuildRuntime(cfg, disp, cwd, env, sysPrompt, opModel, envBlock)
|
rt := session.BuildRuntime(cfg, disp, cwd, env, sysPrompt, opModel, envBlock)
|
||||||
|
|
||||||
restoredSession := agent.RestoreHistory(ps)
|
restoredSession := session.RestoreHistory(ps)
|
||||||
|
|
||||||
var sessPtr *Session
|
var sessPtr *Session
|
||||||
core := agent.New(agent.Config{
|
core := session.New(session.Config{
|
||||||
Backend: be,
|
Backend: be,
|
||||||
AgentName: agentName,
|
AgentName: agentName,
|
||||||
AgentsDir: s.cfg.AgentsDir,
|
AgentsDir: s.cfg.AgentsDir,
|
||||||
|
|
@ -232,7 +232,7 @@ func (s *Manager) restoreSession(ps *agent.PersistedAgent) error {
|
||||||
data := make([]byte, len(sessPtr.plan))
|
data := make([]byte, len(sessPtr.plan))
|
||||||
copy(data, sessPtr.plan)
|
copy(data, sessPtr.plan)
|
||||||
sessPtr.mu.RUnlock()
|
sessPtr.mu.RUnlock()
|
||||||
return agent.NextUncheckedStep(data)
|
return session.NextUncheckedStep(data)
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -1,11 +1,11 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Session holds all state for one agent session.
|
// Session holds all state for one agent session.
|
||||||
|
|
@ -13,7 +13,7 @@ type Session struct {
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
id string
|
id string
|
||||||
uname string // immutable user principal (numeric UID), set at creation
|
uname string // immutable user principal (numeric UID), set at creation
|
||||||
Core *agent.Session
|
Core *session.Session
|
||||||
Ctx context.Context
|
Ctx context.Context
|
||||||
cancel context.CancelFunc
|
cancel context.CancelFunc
|
||||||
log []byte
|
log []byte
|
||||||
|
|
@ -30,7 +30,7 @@ type Session struct {
|
||||||
modelsCacheAt time.Time
|
modelsCacheAt time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewSession(id string, core *agent.Session, ctx context.Context, cancel context.CancelFunc) *Session {
|
func NewSession(id string, core *session.Session, ctx context.Context, cancel context.CancelFunc) *Session {
|
||||||
sess := &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
sess := &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
||||||
sess.startEventLog()
|
sess.startEventLog()
|
||||||
return sess
|
return sess
|
||||||
|
|
@ -83,7 +83,7 @@ func (sess *Session) Cancel() {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sess *Session) Interrupt() {
|
func (sess *Session) Interrupt() {
|
||||||
sess.Core.Interrupt(agent.ErrInterrupted)
|
sess.Core.Interrupt(session.ErrInterrupted)
|
||||||
}
|
}
|
||||||
|
|
||||||
// AppendLog appends data to the session's log and bumps the version.
|
// AppendLog appends data to the session's log and bumps the version.
|
||||||
|
|
@ -113,7 +113,7 @@ func (sess *Session) startEventLog() {
|
||||||
streamingRole := "" // tracks current streaming role ("assistant", "reasoning", or "tool")
|
streamingRole := "" // tracks current streaming role ("assistant", "reasoning", or "tool")
|
||||||
streamingResponseID := ""
|
streamingResponseID := ""
|
||||||
|
|
||||||
sess.Core.Bus().Subscribe("event", func(ev agent.Event) {
|
sess.Core.Bus().Subscribe("event", func(ev session.Event) {
|
||||||
switch ev.Role {
|
switch ev.Role {
|
||||||
case "assistant", "reasoning":
|
case "assistant", "reasoning":
|
||||||
if streamingRole != ev.Role || (ev.Role == "assistant" && ev.ResponseID != streamingResponseID) {
|
if streamingRole != ev.Role || (ev.Role == "assistant" && ev.ResponseID != streamingResponseID) {
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package session
|
package mgr
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
@ -8,7 +8,7 @@ import (
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
agent "ollie/session"
|
"ollie/session"
|
||||||
"ollie/backend"
|
"ollie/backend"
|
||||||
"ollie/tools"
|
"ollie/tools"
|
||||||
olog "ollie/log"
|
olog "ollie/log"
|
||||||
|
|
@ -277,7 +277,7 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) fs.FileSpec {
|
||||||
if base == "" {
|
if base == "" {
|
||||||
base = h.sess.Core.State()
|
base = h.sess.Core.State()
|
||||||
}
|
}
|
||||||
v, ok := h.sess.Core.WaitChange(ctx, agent.WatchState, base)
|
v, ok := h.sess.Core.WaitChange(ctx, session.WatchState, base)
|
||||||
if !ok {
|
if !ok {
|
||||||
s := h.sess.Core.State()
|
s := h.sess.Core.State()
|
||||||
return []byte(s + "\n"), s, nil
|
return []byte(s + "\n"), s, nil
|
||||||
|
|
@ -637,7 +637,7 @@ func (h *sessionHelper) handleCtl(input string) error {
|
||||||
}
|
}
|
||||||
switch cmd[0] {
|
switch cmd[0] {
|
||||||
case "stop":
|
case "stop":
|
||||||
h.sess.Core.Interrupt(agent.ErrInterrupted)
|
h.sess.Core.Interrupt(session.ErrInterrupted)
|
||||||
case "kill":
|
case "kill":
|
||||||
h.kill()
|
h.kill()
|
||||||
case "rn":
|
case "rn":
|
||||||
|
|
@ -21,7 +21,7 @@ import (
|
||||||
olog "ollie/log"
|
olog "ollie/log"
|
||||||
"ollie/paths"
|
"ollie/paths"
|
||||||
"olliesrv/fs"
|
"olliesrv/fs"
|
||||||
"olliesrv/session"
|
"olliesrv/mgr"
|
||||||
|
|
||||||
"9fans.net/go/plan9"
|
"9fans.net/go/plan9"
|
||||||
)
|
)
|
||||||
|
|
@ -68,7 +68,7 @@ type Server struct {
|
||||||
conns []*connState
|
conns []*connState
|
||||||
log *olog.Logger
|
log *olog.Logger
|
||||||
sink *olog.Sink
|
sink *olog.Sink
|
||||||
sessionMgr *session.Manager
|
sessionMgr *mgr.Manager
|
||||||
rootStore *fs.Tree
|
rootStore *fs.Tree
|
||||||
elevateTree *elevateTree
|
elevateTree *elevateTree
|
||||||
groups map[string]map[string]bool // group → set of members
|
groups map[string]map[string]bool // group → set of members
|
||||||
|
|
@ -80,7 +80,7 @@ type Server struct {
|
||||||
// Config holds the pre-built trees and manager for the server.
|
// Config holds the pre-built trees and manager for the server.
|
||||||
type Config struct {
|
type Config struct {
|
||||||
Sink *olog.Sink
|
Sink *olog.Sink
|
||||||
SessionMgr *session.Manager
|
SessionMgr *mgr.Manager
|
||||||
RootStore *fs.Tree
|
RootStore *fs.Tree
|
||||||
ElevateBroker *elevate.Broker
|
ElevateBroker *elevate.Broker
|
||||||
InvalidateModels func()
|
InvalidateModels func()
|
||||||
|
|
@ -491,7 +491,7 @@ func isSessionFile(path string) bool {
|
||||||
if !ok || strings.Contains(name, "/") {
|
if !ok || strings.Contains(name, "/") {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
_, ok = session.FileMode(name)
|
_, ok = mgr.FileMode(name)
|
||||||
return ok
|
return ok
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Reference in New Issue