Remove redundant code and obsolete streaming mechanism
- Extract inject helper for inject/i ctl handlers (dedup) - Remove duplicate metrics file nodes at session/agent level - Remove obsolete StreamFunc/streamWriter from toolsrv exec - Remove stream field from limitedWriter -68 lines
This commit is contained in:
parent
4a037c0bbe
commit
9b9c1a9539
|
|
@ -31,7 +31,6 @@ type (
|
|||
SyntheticFileInfo = virtfs.SyntheticFileInfo
|
||||
)
|
||||
|
||||
|
||||
var treeSpecUID, treeSpecGID string
|
||||
|
||||
const asyncWorkLimit = 64
|
||||
|
|
@ -468,15 +467,6 @@ func buildSessionChildren(
|
|||
return []byte(metrics.Format(a)), nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("metrics", 0444,
|
||||
virtfs.Read(func() ([]byte, error) {
|
||||
a, err := metrics.AggregateSession(s.ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []byte(metrics.Format(a)), nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("metrics.by-backend-model", 0444,
|
||||
virtfs.Read(func() ([]byte, error) {
|
||||
groups, err := metrics.GroupedAggregate(s.ID, "")
|
||||
|
|
@ -726,6 +716,21 @@ func buildAgentChildren(a *agent.Agent, s *session.Session) []virtfs.FsNodeDecl
|
|||
}
|
||||
return []byte(sb.String()), nil
|
||||
}
|
||||
inject := func(args []string) ([]byte, error) {
|
||||
text := strings.Join(args, " ")
|
||||
if text == "" {
|
||||
return nil, fmt.Errorf("inject requires text")
|
||||
}
|
||||
if a.IsRunning() {
|
||||
a.InjectRewrite(text)
|
||||
} else {
|
||||
startAsync(s.Ctx, func() {
|
||||
a.Submit(s.Ctx, text)
|
||||
a.EnsureTrailingNewline()
|
||||
})
|
||||
}
|
||||
return []byte("ok\n"), nil
|
||||
}
|
||||
|
||||
return []virtfs.FsNodeDecl{
|
||||
virtfs.FileNode("prompt", 0666,
|
||||
|
|
@ -961,36 +966,8 @@ func buildAgentChildren(a *agent.Agent, s *session.Session) []virtfs.FsNodeDecl
|
|||
}
|
||||
return []byte("ok\n"), nil
|
||||
},
|
||||
"inject": func(args []string) ([]byte, error) {
|
||||
text := strings.Join(args, " ")
|
||||
if text == "" {
|
||||
return nil, fmt.Errorf("inject requires text")
|
||||
}
|
||||
if a.IsRunning() {
|
||||
a.InjectRewrite(text)
|
||||
} else {
|
||||
startAsync(s.Ctx, func() {
|
||||
a.Submit(s.Ctx, text)
|
||||
a.EnsureTrailingNewline()
|
||||
})
|
||||
}
|
||||
return []byte("ok\n"), nil
|
||||
},
|
||||
"i": func(args []string) ([]byte, error) {
|
||||
text := strings.Join(args, " ")
|
||||
if text == "" {
|
||||
return nil, fmt.Errorf("inject requires text")
|
||||
}
|
||||
if a.IsRunning() {
|
||||
a.InjectRewrite(text)
|
||||
} else {
|
||||
startAsync(s.Ctx, func() {
|
||||
a.Submit(s.Ctx, text)
|
||||
a.EnsureTrailingNewline()
|
||||
})
|
||||
}
|
||||
return []byte("ok\n"), nil
|
||||
},
|
||||
"inject": inject,
|
||||
"i": inject,
|
||||
"agent": func(args []string) ([]byte, error) {
|
||||
if len(args) == 0 {
|
||||
return []byte(a.Profile() + "\n"), nil
|
||||
|
|
@ -1222,15 +1199,6 @@ func buildAgentChildren(a *agent.Agent, s *session.Session) []virtfs.FsNodeDecl
|
|||
return []byte(metrics.Format(a)), nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("metrics", 0444,
|
||||
virtfs.Read(func() ([]byte, error) {
|
||||
a, err := metrics.AggregateAgent(s.ID, a.ID())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []byte(metrics.Format(a)), nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("metrics.by-backend-model", 0444,
|
||||
virtfs.Read(func() ([]byte, error) {
|
||||
groups, err := metrics.GroupedAggregate(s.ID, a.ID())
|
||||
|
|
|
|||
|
|
@ -37,21 +37,6 @@ type Config struct {
|
|||
Started chan StartResult // receives exactly one startup result if set
|
||||
}
|
||||
|
||||
// StreamFunc returns the streaming output function from context, if any.
|
||||
func StreamFunc(ctx context.Context) func(string) {
|
||||
if fn, ok := ctx.Value(streamKey{}).(func(string)); ok {
|
||||
return fn
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type streamKey struct{}
|
||||
|
||||
// WithStreamFunc attaches a streaming output function to context.
|
||||
func WithStreamFunc(ctx context.Context, fn func(string)) context.Context {
|
||||
return context.WithValue(ctx, streamKey{}, fn)
|
||||
}
|
||||
|
||||
// ExecuteTool runs a tool script with the given args and returns the result.
|
||||
func ExecuteTool(ctx context.Context, info protocol.ToolInfo, args json.RawMessage, cfg Config) (json.RawMessage, error) {
|
||||
// Resolve script path
|
||||
|
|
@ -208,9 +193,8 @@ func executeSandboxed(ctx context.Context, toolPath, stdinData, cwd string, envE
|
|||
w = io.MultiWriter(&outputBuf, streamOut)
|
||||
}
|
||||
lw := &limitedWriter{
|
||||
w: w,
|
||||
limit: 10 * 1024 * 1024,
|
||||
stream: StreamFunc(ctx),
|
||||
w: w,
|
||||
limit: 10 * 1024 * 1024,
|
||||
}
|
||||
cmd.Stdout = lw
|
||||
cmd.Stderr = lw
|
||||
|
|
@ -302,12 +286,6 @@ func executeDirectUnsandboxed(ctx context.Context, cmd, cwd string, env map[stri
|
|||
execCmd.Stdout = &stdout
|
||||
execCmd.Stderr = &stderr
|
||||
|
||||
// Stream output if context has a stream function
|
||||
if stream := StreamFunc(ctx); stream != nil {
|
||||
execCmd.Stdout = io.MultiWriter(&stdout, &streamWriter{stream})
|
||||
execCmd.Stderr = io.MultiWriter(&stderr, &streamWriter{stream})
|
||||
}
|
||||
|
||||
if err := execCmd.Start(); err != nil {
|
||||
if started != nil {
|
||||
started <- StartResult{Err: fmt.Errorf("bypass execution failed: %w", err)}
|
||||
|
|
@ -335,16 +313,6 @@ func executeDirectUnsandboxed(ctx context.Context, cmd, cwd string, env map[stri
|
|||
return output, nil
|
||||
}
|
||||
|
||||
// streamWriter adapts a stream function to io.Writer.
|
||||
type streamWriter struct {
|
||||
fn func(string)
|
||||
}
|
||||
|
||||
func (w *streamWriter) Write(p []byte) (int, error) {
|
||||
w.fn(string(p))
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
// Plan9Namespace computes the Plan 9 namespace directory.
|
||||
func Plan9Namespace(env map[string]string) string {
|
||||
if ns := env["NAMESPACE"]; ns != "" {
|
||||
|
|
@ -360,13 +328,12 @@ func Plan9Namespace(env map[string]string) string {
|
|||
return "/tmp/ns." + env["USER"] + "." + disp
|
||||
}
|
||||
|
||||
// limitedWriter wraps a writer with a size limit and optional streaming.
|
||||
// limitedWriter wraps a writer with a size limit.
|
||||
type limitedWriter struct {
|
||||
w io.Writer
|
||||
limit int
|
||||
written int
|
||||
truncated bool
|
||||
stream func(string)
|
||||
}
|
||||
|
||||
func (lw *limitedWriter) Write(p []byte) (int, error) {
|
||||
|
|
@ -385,8 +352,5 @@ func (lw *limitedWriter) Write(p []byte) (int, error) {
|
|||
}
|
||||
n, err := lw.w.Write(toWrite)
|
||||
lw.written += n
|
||||
if lw.stream != nil && n > 0 {
|
||||
lw.stream(string(toWrite[:n]))
|
||||
}
|
||||
return len(p), err
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue