all: merge toolsrv files, fix context hierarchy, remove dead params
toolsrv: - Merge rpc.go into conn.go (wire types already there) - Merge tools.go + discover.go into registry.go - Merge accessors.go into server.go - 4 files eliminated fs: - HandlerCtx.Context now carries the session context (set via sessionBindings Applier). Previously it was the daemon context which never cancels — handlers checking ctx.Done() now properly detect session kill/pause. - Session-scoped blocking handlers (streamAgentChat, blockAgentStateWait) merge session context with per-read timeout via mergeCtx(). A session kill now immediately unblocks readers without waiting for the 5s timeout. - Remove redundant manual AfterFunc in blockAgentStateWait. session: - Remove unused _ string parameter from WaitEvent
This commit is contained in:
parent
5b64d62d2e
commit
455434608a
|
|
@ -18,6 +18,14 @@ import (
|
|||
"ollie/toolsrv"
|
||||
)
|
||||
|
||||
// mergeCtx returns a context that cancels when either parent or child cancels.
|
||||
// parent is typically the session context; child is the per-read timeout from the server.
|
||||
func mergeCtx(parent, child context.Context) (context.Context, context.CancelFunc) {
|
||||
merged, cancel := context.WithCancel(parent)
|
||||
stop := context.AfterFunc(child, func() { cancel() })
|
||||
return merged, func() { stop(); cancel() }
|
||||
}
|
||||
|
||||
// wireAgentStateEvents sets up the agent to publish state change events to the event bus.
|
||||
func wireAgentStateEvents(sessID string, ag *agent.Agent) {
|
||||
ag.SetOnStateChange(func(agentID, state string) {
|
||||
|
|
@ -113,8 +121,8 @@ func readEventwait(_ HandlerCtx) ([]byte, error) {
|
|||
return nil, fmt.Errorf("eventwait: use blocking read")
|
||||
}
|
||||
|
||||
func blockEventwait(_ HandlerCtx, cctx context.Context, base string) ([]byte, string, error) {
|
||||
return session.WaitEvent(cctx, base)
|
||||
func blockEventwait(_ HandlerCtx, cctx context.Context, _ string) ([]byte, string, error) {
|
||||
return session.WaitEvent(cctx)
|
||||
}
|
||||
|
||||
func requestGenerate(ctx HandlerCtx, data []byte) ([]byte, error) {
|
||||
|
|
@ -253,6 +261,7 @@ func sessionBindings(ctx HandlerCtx) ([]Binding, error) {
|
|||
UID: sess.Uname(),
|
||||
GID: "agent",
|
||||
Applier: func(c HandlerCtx) HandlerCtx {
|
||||
c.Context = node.Ctx()
|
||||
c.Session = node
|
||||
c.Remove = func() {
|
||||
session.Kill(n)
|
||||
|
|
@ -485,7 +494,9 @@ func readAgentChat(ctx HandlerCtx) ([]byte, error) {
|
|||
}
|
||||
|
||||
func streamAgentChat(ctx HandlerCtx, cctx context.Context, base string) ([]byte, string, error) {
|
||||
return streamChat(ctx.AgentLog, cctx, base)
|
||||
merged, cancel := mergeCtx(ctx, cctx)
|
||||
defer cancel()
|
||||
return streamChat(ctx.AgentLog, merged, base)
|
||||
}
|
||||
|
||||
func readAgentConnection(ctx HandlerCtx) ([]byte, error) {
|
||||
|
|
@ -499,13 +510,12 @@ func readAgentConnection(ctx HandlerCtx) ([]byte, error) {
|
|||
}
|
||||
|
||||
func blockAgentStateWait(ctx HandlerCtx, cctx context.Context, base string) ([]byte, string, error) {
|
||||
waitCtx, cancel := context.WithCancel(cctx)
|
||||
merged, cancel := mergeCtx(ctx, cctx)
|
||||
defer cancel()
|
||||
context.AfterFunc(ctx.Session.Ctx(), cancel)
|
||||
if base == "" {
|
||||
base = ctx.Agent.State()
|
||||
}
|
||||
v, ok := ctx.Agent.WaitChange(waitCtx, agent.WatchState, base)
|
||||
v, ok := ctx.Agent.WaitChange(merged, agent.WatchState, base)
|
||||
if !ok {
|
||||
st := ctx.Agent.State()
|
||||
return []byte(st + "\n"), st, nil
|
||||
|
|
|
|||
|
|
@ -56,7 +56,7 @@ func SubscribeEvents(ctx context.Context, topic string) <-chan Event {
|
|||
}
|
||||
|
||||
// WaitEvent blocks until an event is available. Returns the event as "topic payload\n".
|
||||
func WaitEvent(ctx context.Context, _ string) ([]byte, string, error) {
|
||||
func WaitEvent(ctx context.Context) ([]byte, string, error) {
|
||||
ch := SubscribeEvents(ctx, "*")
|
||||
select {
|
||||
case ev, ok := <-ch:
|
||||
|
|
|
|||
|
|
@ -1,33 +0,0 @@
|
|||
package toolsrv
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// SetOnToolsChanged sets the callback for tool directory changes.
|
||||
func (s *Server) SetOnToolsChanged(fn func(string)) {
|
||||
s.OnToolsChanged = fn
|
||||
}
|
||||
|
||||
// ToolRegistryRevision returns the current registry revision for this server's session.
|
||||
// Returns 0 if no registry is configured.
|
||||
func (s *Server) ToolRegistryRevision() uint64 {
|
||||
if s.toolRegistry == nil || s.sessionID == "" {
|
||||
return 0
|
||||
}
|
||||
return s.toolRegistry.Revision(s.sessionID)
|
||||
}
|
||||
|
||||
// buildToolListing formats a slice of ToolInfo into the preamble listing
|
||||
// format consumed by agent.refreshToolListing: "- **name** — description\n" per entry.
|
||||
// Only local tools (Server == "") with a non-empty description are included.
|
||||
func buildToolListing(infos []ToolInfo) string {
|
||||
var sb strings.Builder
|
||||
for _, ti := range infos {
|
||||
if ti.Description != "" && ti.Server == "" {
|
||||
fmt.Fprintf(&sb, "- **%s** — %s\n", ti.Name, ti.Description)
|
||||
}
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
|
|
@ -281,3 +281,29 @@ func (c *Conn) callStreaming(ctx context.Context, method string, params json.Raw
|
|||
return nil, fmt.Errorf("rpc %s: %w", method, err)
|
||||
}
|
||||
}
|
||||
|
||||
// JSON-RPC 2.0 types shared between Conn (client) and transport layer.
|
||||
|
||||
type rpcRequest struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int64 `json:"id"`
|
||||
Method string `json:"method"`
|
||||
Params json.RawMessage `json:"params,omitempty"`
|
||||
}
|
||||
|
||||
type rpcResponse struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int64 `json:"id"`
|
||||
Result json.RawMessage `json:"result,omitempty"`
|
||||
Error *rpcError `json:"error,omitempty"`
|
||||
Stream bool `json:"stream,omitempty"`
|
||||
}
|
||||
|
||||
type rpcError struct {
|
||||
Code int `json:"code"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
type outputNotification struct {
|
||||
Data string `json:"data"`
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,46 +0,0 @@
|
|||
package toolsrv
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"ollie/paths"
|
||||
)
|
||||
|
||||
// ToolsPath returns ~/.config/ollie/tools.
|
||||
func ToolsPath() string {
|
||||
return paths.CfgDir() + "/tools"
|
||||
}
|
||||
|
||||
// DiscoverTools scans the tools directory for .meta files and returns
|
||||
// metadata for all tools.
|
||||
func DiscoverTools() []ToolInfo {
|
||||
dir := ToolsPath()
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
var infos []ToolInfo
|
||||
for _, e := range entries {
|
||||
if e.IsDir() || !strings.HasSuffix(e.Name(), ".meta") {
|
||||
continue
|
||||
}
|
||||
name := strings.TrimSuffix(e.Name(), ".meta")
|
||||
data, err := os.ReadFile(filepath.Join(dir, e.Name()))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
var m MetaFile
|
||||
if err := json.Unmarshal(data, &m); err != nil {
|
||||
continue
|
||||
}
|
||||
info := ToolInfoFromMeta(name, &m)
|
||||
if info.Name == "" {
|
||||
continue // no variant matched — tool unavailable on this host
|
||||
}
|
||||
infos = append(infos, info)
|
||||
}
|
||||
return infos
|
||||
}
|
||||
|
|
@ -1,10 +1,15 @@
|
|||
package toolsrv
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"ollie/paths"
|
||||
)
|
||||
|
||||
type Registry struct {
|
||||
|
|
@ -110,3 +115,62 @@ func (r *Registry) Revision(sessionID string) uint64 {
|
|||
}
|
||||
|
||||
|
||||
|
||||
// ToolInfo describes a tool provided by a server.
|
||||
type ToolInfo struct {
|
||||
Server string
|
||||
Name string
|
||||
Description string
|
||||
InputSchema json.RawMessage
|
||||
// Prompt is the usage documentation for this tool, extracted from
|
||||
// the script's ollie:prompt block. Included in the system prompt.
|
||||
Prompt string
|
||||
// Tier is the retention tier: "hot", "warm", or "cold". Parsed from
|
||||
// the script's ollie:tier annotation. Empty defaults to "hot".
|
||||
Tier string
|
||||
// ReadOnly is true when the tool carries an "ollie:parallel read" annotation.
|
||||
ReadOnly bool
|
||||
// OutputFormat is the source-fence language for tool output. Empty means plaintext.
|
||||
OutputFormat string
|
||||
// ResetsCounter is true for action-oriented tools (file_write, lsp_rename, etc.)
|
||||
// that indicate active progress rather than passive research. When true, the
|
||||
// agent loop resets its step counter, allowing continued work without hitting
|
||||
// the soft step-budget guardrail.
|
||||
ResetsCounter bool
|
||||
}
|
||||
|
||||
// ToolsPath returns ~/.config/ollie/tools.
|
||||
func ToolsPath() string {
|
||||
return paths.CfgDir() + "/tools"
|
||||
}
|
||||
|
||||
// DiscoverTools scans the tools directory for .meta files and returns
|
||||
// metadata for all tools.
|
||||
func DiscoverTools() []ToolInfo {
|
||||
dir := ToolsPath()
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
var infos []ToolInfo
|
||||
for _, e := range entries {
|
||||
if e.IsDir() || !strings.HasSuffix(e.Name(), ".meta") {
|
||||
continue
|
||||
}
|
||||
name := strings.TrimSuffix(e.Name(), ".meta")
|
||||
data, err := os.ReadFile(filepath.Join(dir, e.Name()))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
var m MetaFile
|
||||
if err := json.Unmarshal(data, &m); err != nil {
|
||||
continue
|
||||
}
|
||||
info := ToolInfoFromMeta(name, &m)
|
||||
if info.Name == "" {
|
||||
continue // no variant matched — tool unavailable on this host
|
||||
}
|
||||
infos = append(infos, info)
|
||||
}
|
||||
return infos
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,29 +0,0 @@
|
|||
package toolsrv
|
||||
|
||||
import "encoding/json"
|
||||
|
||||
// JSON-RPC 2.0 types shared between Conn (client) and transport layer.
|
||||
|
||||
type rpcRequest struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int64 `json:"id"`
|
||||
Method string `json:"method"`
|
||||
Params json.RawMessage `json:"params,omitempty"`
|
||||
}
|
||||
|
||||
type rpcResponse struct {
|
||||
JSONRPC string `json:"jsonrpc"`
|
||||
ID int64 `json:"id"`
|
||||
Result json.RawMessage `json:"result,omitempty"`
|
||||
Error *rpcError `json:"error,omitempty"`
|
||||
Stream bool `json:"stream,omitempty"`
|
||||
}
|
||||
|
||||
type rpcError struct {
|
||||
Code int `json:"code"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
type outputNotification struct {
|
||||
Data string `json:"data"`
|
||||
}
|
||||
|
|
@ -392,3 +392,30 @@ func (s *Server) cleanupDetached() {
|
|||
p.Mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// SetOnToolsChanged sets the callback for tool directory changes.
|
||||
func (s *Server) SetOnToolsChanged(fn func(string)) {
|
||||
s.OnToolsChanged = fn
|
||||
}
|
||||
|
||||
// ToolRegistryRevision returns the current registry revision for this server's session.
|
||||
// Returns 0 if no registry is configured.
|
||||
func (s *Server) ToolRegistryRevision() uint64 {
|
||||
if s.toolRegistry == nil || s.sessionID == "" {
|
||||
return 0
|
||||
}
|
||||
return s.toolRegistry.Revision(s.sessionID)
|
||||
}
|
||||
|
||||
// buildToolListing formats a slice of ToolInfo into the preamble listing
|
||||
// format consumed by agent.refreshToolListing: "- **name** — description\n" per entry.
|
||||
// Only local tools (Server == "") with a non-empty description are included.
|
||||
func buildToolListing(infos []ToolInfo) string {
|
||||
var sb strings.Builder
|
||||
for _, ti := range infos {
|
||||
if ti.Description != "" && ti.Server == "" {
|
||||
fmt.Fprintf(&sb, "- **%s** — %s\n", ti.Name, ti.Description)
|
||||
}
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,30 +0,0 @@
|
|||
// Package tools implements the tool server (sandboxed execution, tool registry,
|
||||
// skill management) and supporting types.
|
||||
package toolsrv
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
)
|
||||
|
||||
// ToolInfo describes a tool provided by a server.
|
||||
type ToolInfo struct {
|
||||
Server string
|
||||
Name string
|
||||
Description string
|
||||
InputSchema json.RawMessage
|
||||
// Prompt is the usage documentation for this tool, extracted from
|
||||
// the script's ollie:prompt block. Included in the system prompt.
|
||||
Prompt string
|
||||
// Tier is the retention tier: "hot", "warm", or "cold". Parsed from
|
||||
// the script's ollie:tier annotation. Empty defaults to "hot".
|
||||
Tier string
|
||||
// ReadOnly is true when the tool carries an "ollie:parallel read" annotation.
|
||||
ReadOnly bool
|
||||
// OutputFormat is the source-fence language for tool output. Empty means plaintext.
|
||||
OutputFormat string
|
||||
// ResetsCounter is true for action-oriented tools (file_write, lsp_rename, etc.)
|
||||
// that indicate active progress rather than passive research. When true, the
|
||||
// agent loop resets its step counter, allowing continued work without hitting
|
||||
// the soft step-budget guardrail.
|
||||
ResetsCounter bool
|
||||
}
|
||||
Loading…
Reference in New Issue