refactor: extract format package, ctl dispatcher, and log structure validator
- Extract all [[[block]]] formatting constants/helpers into format/format.go (EndTag, RoleDelim, ToolDelim, CallDelim, BlockHeader, ParseBlockHeader, IsEndMarker) - Replace switch/case ctl dispatchers in fs/handlers.go with generic ctlDispatch() from new fs/ctl.go - Add LogValidator state machine to AgentLog enforcing: - Every [[[block]]] MUST have an [[[end]]] marker - Blocks SHALL NOT be nested - Streaming line buffering across partial writes - Fatal failure on first violation with line-number reporting - UnclosedAtTeardown() for session/agent teardown checks - Fix kde/kate ollie_ghost compile: stripPrefixEcho signature mismatch
This commit is contained in:
parent
7c99d101b1
commit
c7f1bec71c
|
|
@ -0,0 +1,87 @@
|
|||
package format
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Formatting constants for the 9P chat log protocol.
|
||||
// Used by fs.NewEventHandler and fs.replayMessagesToLog.
|
||||
|
||||
const (
|
||||
// EndTag marks the end of a message block.
|
||||
EndTag = "[[[end]]]"
|
||||
|
||||
// RoleTag formats a role delimiter, e.g. "[[[user]]]".
|
||||
RoleTag = "[[[%s]]]"
|
||||
|
||||
// ToolFence is the opening fence for tool output blocks.
|
||||
ToolFence = "```\n"
|
||||
|
||||
// AssistantHeaderPrefix is prepended to assistant headers when a response ID is present.
|
||||
AssistantHeaderPrefix = "[[[assistant:"
|
||||
|
||||
// AssistantTag is the bare assistant header (no response ID).
|
||||
AssistantTag = "[[[assistant]]]"
|
||||
|
||||
// UserTag is the user role delimiter.
|
||||
UserTag = "[[[user]]]"
|
||||
|
||||
// RetryTag is the retry role delimiter.
|
||||
RetryTag = "[[[retry]]]"
|
||||
|
||||
// ToolDelimPrefix is the prefix for tool delimiters, e.g. "[[[tool:shell]]]".
|
||||
ToolDelimPrefix = "[[[tool:"
|
||||
|
||||
// ToolTag is the bare tool role delimiter (no name).
|
||||
ToolTag = "[[[tool]]]"
|
||||
|
||||
// CallDelimPrefix is the prefix for call delimiters, e.g. "[[[call:file_read]]]".
|
||||
CallDelimPrefix = "[[[call:"
|
||||
)
|
||||
|
||||
// RoleDelim returns a formatted role delimiter.
|
||||
func RoleDelim(role string) string {
|
||||
return fmt.Sprintf(RoleTag, role)
|
||||
}
|
||||
|
||||
// ToolDelim returns a formatted tool delimiter.
|
||||
func ToolDelim(name string) string {
|
||||
return fmt.Sprintf(ToolDelimPrefix+"%s]]]\n", name)
|
||||
}
|
||||
|
||||
// CallDelim returns a formatted call delimiter.
|
||||
func CallDelim(name string) string {
|
||||
return fmt.Sprintf(CallDelimPrefix+"%s]]]\n", name)
|
||||
}
|
||||
|
||||
// BlockHeader is a single-line [[[...]]] marker that opens a log block.
|
||||
type BlockHeader struct {
|
||||
Role string // e.g. "user", "assistant", "tool", "retry", or "tool:shell"
|
||||
Name string // for tool/call: the sub-name (e.g. "shell"), empty for bare roles
|
||||
}
|
||||
|
||||
// ParseBlockHeader extracts a BlockHeader from a line like "[[[tool:shell]]]"
|
||||
// or "[[[assistant]]]". Returns nil if the line is not a valid block header.
|
||||
func ParseBlockHeader(line string) *BlockHeader {
|
||||
const prefix, suffix = "[[[", "]]]"
|
||||
if !strings.HasPrefix(line, prefix) || !strings.HasSuffix(line, suffix) {
|
||||
return nil
|
||||
}
|
||||
inner := strings.TrimPrefix(line, prefix)
|
||||
inner = strings.TrimSuffix(inner, suffix)
|
||||
if inner == "" {
|
||||
return nil
|
||||
}
|
||||
bh := &BlockHeader{Role: inner}
|
||||
if idx := strings.Index(inner, ":"); idx >= 0 {
|
||||
bh.Role = inner[:idx]
|
||||
bh.Name = inner[idx+1:]
|
||||
}
|
||||
return bh
|
||||
}
|
||||
|
||||
// IsEndMarker reports whether the line is exactly the [[[end]]] closing marker.
|
||||
func IsEndMarker(line string) bool {
|
||||
return line == EndTag
|
||||
}
|
||||
|
|
@ -0,0 +1,28 @@
|
|||
package fs
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// ctlHandler is a function that handles a single ctl command.
|
||||
type ctlHandler func(args []string, ctx HandlerCtx) error
|
||||
|
||||
// ctlDispatch returns a Write handler that dispatches based on the first
|
||||
// whitespace-delimited word in the input. The map keys are command names;
|
||||
// values are the handlers. If the command is unknown, an error is returned.
|
||||
func ctlDispatch(handlers map[string]ctlHandler) func(HandlerCtx, []byte) error {
|
||||
return func(ctx HandlerCtx, data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
parts := strings.Fields(input)
|
||||
cmd := parts[0]
|
||||
h, ok := handlers[cmd]
|
||||
if !ok {
|
||||
return fmt.Errorf("unknown ctl command: %s", cmd)
|
||||
}
|
||||
return h(parts[1:], ctx)
|
||||
}
|
||||
}
|
||||
159
fs/handlers.go
159
fs/handlers.go
|
|
@ -89,24 +89,24 @@ func readTools(_ HandlerCtx) ([]byte, error) {
|
|||
}
|
||||
|
||||
func writeRootCtl(ctx HandlerCtx, data []byte) error {
|
||||
cmd := strings.TrimSpace(string(data))
|
||||
switch cmd {
|
||||
case "invalidate":
|
||||
if ctx.Models != nil {
|
||||
ctx.Models.Invalidate()
|
||||
}
|
||||
if ctx.Invalidate != nil {
|
||||
ctx.Invalidate()
|
||||
}
|
||||
return nil
|
||||
case "kill":
|
||||
if ctx.Shutdown != nil {
|
||||
ctx.Shutdown()
|
||||
}
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("unknown ctl command: %s (valid: invalidate, kill)", cmd)
|
||||
handlers := map[string]ctlHandler{
|
||||
"invalidate": func(_ []string, _ HandlerCtx) error {
|
||||
if ctx.Models != nil {
|
||||
ctx.Models.Invalidate()
|
||||
}
|
||||
if ctx.Invalidate != nil {
|
||||
ctx.Invalidate()
|
||||
}
|
||||
return nil
|
||||
},
|
||||
"kill": func(_ []string, _ HandlerCtx) error {
|
||||
if ctx.Shutdown != nil {
|
||||
ctx.Shutdown()
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
return ctlDispatch(handlers)(ctx, data)
|
||||
}
|
||||
|
||||
func readEventwait(_ HandlerCtx) ([]byte, error) {
|
||||
|
|
@ -311,26 +311,35 @@ func readSessionConnected(ctx HandlerCtx) ([]byte, error) {
|
|||
}
|
||||
|
||||
func writeSessionCtl(ctx HandlerCtx, data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
handlers := map[string]ctlHandler{
|
||||
"kill": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Remove()
|
||||
return nil
|
||||
},
|
||||
".": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Remove()
|
||||
return nil
|
||||
},
|
||||
"save": func(_ []string, _ HandlerCtx) error {
|
||||
if ctx.Session.Core != nil {
|
||||
ctx.Session.Core.SaveSession("")
|
||||
}
|
||||
return nil
|
||||
},
|
||||
"invalidate": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Session.InvalidateModelsCache()
|
||||
return nil
|
||||
},
|
||||
"pause": func(_ []string, _ HandlerCtx) error {
|
||||
return ctx.Session.Pause()
|
||||
},
|
||||
"resume": func(_ []string, _ HandlerCtx) error {
|
||||
return ctx.Session.Resume()
|
||||
},
|
||||
}
|
||||
cmd := strings.Fields(input)
|
||||
switch cmd[0] {
|
||||
case "kill", ".":
|
||||
ctx.Remove()
|
||||
case "save":
|
||||
if ctx.Session.Core != nil {
|
||||
ctx.Session.Core.SaveSession("")
|
||||
}
|
||||
case "invalidate":
|
||||
ctx.Session.InvalidateModelsCache()
|
||||
case "pause":
|
||||
return ctx.Session.Pause()
|
||||
case "resume":
|
||||
return ctx.Session.Resume()
|
||||
default:
|
||||
return fmt.Errorf("unknown session ctl command: %s", cmd[0])
|
||||
err := ctlDispatch(handlers)(ctx, data)
|
||||
if err != nil {
|
||||
return fmt.Errorf("session ctl: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
@ -576,42 +585,50 @@ func writeAgentCfg(ctx HandlerCtx, data []byte) error {
|
|||
}
|
||||
|
||||
func writeAgentCtl(ctx HandlerCtx, data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
cmd := strings.Fields(input)
|
||||
if len(cmd) == 0 {
|
||||
return nil
|
||||
}
|
||||
switch cmd[0] {
|
||||
case "kill":
|
||||
if ctx.Session.Core != nil {
|
||||
ctx.Session.Core.RemoveAgent(ctx.Agent.ID())
|
||||
session.PublishEvent("session."+ctx.Session.ID()+".agent."+ctx.Agent.ID()+".kill", "")
|
||||
}
|
||||
case "stop":
|
||||
ctx.Agent.Interrupt(agent.ErrInterrupted)
|
||||
case "compact":
|
||||
ctx.Agent.HandleCommand(ctx.Session.Ctx(), "/compact")
|
||||
case "clear":
|
||||
ctx.Agent.HandleCommand(ctx.Session.Ctx(), "/clear")
|
||||
case "model":
|
||||
if len(cmd) > 1 {
|
||||
if be := ctx.Agent.Backend(); be != nil {
|
||||
be.SetModel(strings.Join(cmd[1:], " "))
|
||||
handlers := map[string]ctlHandler{
|
||||
"kill": func(_ []string, _ HandlerCtx) error {
|
||||
if ctx.Session.Core != nil {
|
||||
ctx.Session.Core.RemoveAgent(ctx.Agent.ID())
|
||||
session.PublishEvent("session."+ctx.Session.ID()+".agent."+ctx.Agent.ID()+".kill", "")
|
||||
}
|
||||
}
|
||||
case "name":
|
||||
if len(cmd) > 1 {
|
||||
ctx.Agent.SetName(strings.Join(cmd[1:], " "))
|
||||
}
|
||||
case "cwd":
|
||||
if len(cmd) > 1 {
|
||||
ctx.Agent.SetCWD(strings.Join(cmd[1:], " "))
|
||||
}
|
||||
default:
|
||||
return fmt.Errorf("unknown agent ctl command: %s", cmd[0])
|
||||
return nil
|
||||
},
|
||||
"stop": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Agent.Interrupt(agent.ErrInterrupted)
|
||||
return nil
|
||||
},
|
||||
"compact": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Agent.HandleCommand(ctx.Session.Ctx(), "/compact")
|
||||
return nil
|
||||
},
|
||||
"clear": func(_ []string, _ HandlerCtx) error {
|
||||
ctx.Agent.HandleCommand(ctx.Session.Ctx(), "/clear")
|
||||
return nil
|
||||
},
|
||||
"model": func(args []string, _ HandlerCtx) error {
|
||||
if len(args) > 0 {
|
||||
if be := ctx.Agent.Backend(); be != nil {
|
||||
be.SetModel(strings.Join(args, " "))
|
||||
}
|
||||
}
|
||||
return nil
|
||||
},
|
||||
"name": func(args []string, _ HandlerCtx) error {
|
||||
if len(args) > 0 {
|
||||
ctx.Agent.SetName(strings.Join(args, " "))
|
||||
}
|
||||
return nil
|
||||
},
|
||||
"cwd": func(args []string, _ HandlerCtx) error {
|
||||
if len(args) > 0 {
|
||||
ctx.Agent.SetCWD(paths.ExpandHome(strings.Join(args, " ")))
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
err := ctlDispatch(handlers)(ctx, data)
|
||||
if err != nil {
|
||||
return fmt.Errorf("agent ctl: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,153 @@
|
|||
package fs
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"ollie/format"
|
||||
)
|
||||
|
||||
// block tracks a single open [[[block]]] for nesting validation.
|
||||
type block struct {
|
||||
header format.BlockHeader
|
||||
line int // 1-indexed line number where this block opened
|
||||
}
|
||||
|
||||
// LogValidator enforces mechanical guarantees on the [[[block]]]/[[[end]]] protocol:
|
||||
// - Every [[[block]]] MUST have an [[[end]]] marker
|
||||
// - Blocks SHALL NOT be nested
|
||||
//
|
||||
// It operates as a streaming state machine, buffering partial lines across calls
|
||||
// to Validate() so that block headers split across writes are handled correctly.
|
||||
type LogValidator struct {
|
||||
mu sync.Mutex
|
||||
stack []block // stack of currently open blocks (should never exceed depth 1)
|
||||
buf string // incomplete last line accumulated across calls
|
||||
errs []violation // accumulated validation errors
|
||||
fatal bool // once set, all subsequent data is rejected
|
||||
lineNum int // 1-indexed line counter within buf
|
||||
}
|
||||
|
||||
// violation records a single structural protocol violation.
|
||||
type violation struct {
|
||||
line int
|
||||
msg string
|
||||
}
|
||||
|
||||
func (v violation) Error() string {
|
||||
return fmt.Sprintf("log structure violation at line %d: %s", v.line, v.msg)
|
||||
}
|
||||
|
||||
// NewLogValidator creates a fresh validator with no open blocks.
|
||||
func NewLogValidator() *LogValidator {
|
||||
return &LogValidator{}
|
||||
}
|
||||
|
||||
// Validate processes a chunk of log data through the state machine.
|
||||
// It buffers incomplete lines across calls so that block markers split
|
||||
// across writes are handled correctly. Returns nil if valid, or a slice
|
||||
// of violations if any structural rules are broken.
|
||||
//
|
||||
// The state machine recognizes two kinds of lines:
|
||||
// - Block headers: lines matching "[[[<role>]]]" or "[[[<role>:<name>]]]"
|
||||
// - End markers: the exact line "[[[end]]]"
|
||||
//
|
||||
// Rules enforced:
|
||||
// 1. A block header while another block is open → "nested block" error
|
||||
// 2. An end marker with no open block → "unexpected end" error
|
||||
// 3. At session end, unclosed blocks are reported via UnclosedBlocks()
|
||||
func (v *LogValidator) Validate(data []byte) []error {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
if v.fatal {
|
||||
return []error{fmt.Errorf("validator already failed — rejecting further data")}
|
||||
}
|
||||
|
||||
v.buf += string(data)
|
||||
|
||||
lines := strings.Split(v.buf, "\n")
|
||||
v.buf = lines[len(lines)-1] // keep incomplete trailing content
|
||||
|
||||
for i := 0; i < len(lines)-1; i++ {
|
||||
line := strings.TrimSpace(lines[i])
|
||||
if line == "" {
|
||||
continue
|
||||
}
|
||||
v.lineNum++
|
||||
|
||||
if format.IsEndMarker(line) {
|
||||
v.handleEndMarker(line)
|
||||
} else if bh := format.ParseBlockHeader(line); bh != nil {
|
||||
v.handleBlockHeader(bh, line)
|
||||
}
|
||||
// Lines that are neither block headers nor end markers are ignored
|
||||
// (they're block body content).
|
||||
}
|
||||
|
||||
if v.fatal {
|
||||
errs := make([]error, len(v.errs))
|
||||
for i, e := range v.errs {
|
||||
errs[i] = e
|
||||
}
|
||||
return errs
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (v *LogValidator) handleEndMarker(line string) {
|
||||
if len(v.stack) == 0 {
|
||||
v.errs = append(v.errs, violation{
|
||||
line: v.lineNum,
|
||||
msg: fmt.Sprintf("unexpected [%s] with no open block", line),
|
||||
})
|
||||
v.fatal = true
|
||||
return
|
||||
}
|
||||
top := v.stack[len(v.stack)-1]
|
||||
v.stack = v.stack[:len(v.stack)-1]
|
||||
_ = top // consumed
|
||||
}
|
||||
|
||||
func (v *LogValidator) handleBlockHeader(bh *format.BlockHeader, line string) {
|
||||
if len(v.stack) > 0 {
|
||||
outer := v.stack[len(v.stack)-1]
|
||||
v.errs = append(v.errs, violation{
|
||||
line: v.lineNum,
|
||||
msg: fmt.Sprintf("nested block [%s] while [%s:%s] is still open (no nesting allowed)",
|
||||
line, outer.header.Role, outer.header.Name),
|
||||
})
|
||||
v.fatal = true
|
||||
return
|
||||
}
|
||||
v.stack = append(v.stack, block{header: *bh, line: v.lineNum})
|
||||
}
|
||||
|
||||
// UnclosedBlocks returns a list of blocks that were opened but never closed.
|
||||
// Call this at session/agent teardown to detect missing [[[end]]] markers.
|
||||
func (v *LogValidator) UnclosedBlocks() []block {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
out := make([]block, len(v.stack))
|
||||
copy(out, v.stack)
|
||||
return out
|
||||
}
|
||||
|
||||
// HasViolations reports whether any violations have been recorded.
|
||||
func (v *LogValidator) HasViolations() bool {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
return v.fatal || len(v.errs) > 0
|
||||
}
|
||||
|
||||
// Errors returns all accumulated violations (non-nil only after fatal failure).
|
||||
func (v *LogValidator) Errors() []error {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
out := make([]error, len(v.errs))
|
||||
for i, e := range v.errs {
|
||||
out[i] = e
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
|
@ -0,0 +1,168 @@
|
|||
package fs
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"ollie/format"
|
||||
)
|
||||
|
||||
func TestValidate_validUserAssistant(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[user]]]\nhello\n[[[end]]]\n[[[assistant]]]\nworld\n[[[end]]]\n")
|
||||
errs := v.Validate(data)
|
||||
if errs != nil {
|
||||
t.Fatalf("expected no errors, got: %v", errs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_nestedBlocks(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[user]]]\nouter\n[[[assistant]]]\nnested — should fail\n")
|
||||
errs := v.Validate(data)
|
||||
if len(errs) == 0 {
|
||||
t.Fatal("expected nesting violation, got none")
|
||||
}
|
||||
got := errs[0].Error()
|
||||
if !strings.Contains(got, "nested block") {
|
||||
t.Errorf("expected 'nested block' in error, got: %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_unexpectedEnd(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[end]]]\n")
|
||||
errs := v.Validate(data)
|
||||
if len(errs) == 0 {
|
||||
t.Fatal("expected unexpected end error, got none")
|
||||
}
|
||||
got := errs[0].Error()
|
||||
if !strings.Contains(got, "unexpected") {
|
||||
t.Errorf("expected 'unexpected' in error, got: %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_unclosedAtTeardown(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[user]]]\nhello\n")
|
||||
v.Validate(data)
|
||||
unclosed := v.UnclosedBlocks()
|
||||
if len(unclosed) != 1 {
|
||||
t.Fatalf("expected 1 unclosed block, got %d", len(unclosed))
|
||||
}
|
||||
if unclosed[0].header.Role != "user" {
|
||||
t.Errorf("expected role 'user', got '%s'", unclosed[0].header.Role)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_toolWithSubname(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[tool:shell]]]\necho hi\n```\n[[[end]]]\n")
|
||||
errs := v.Validate(data)
|
||||
if errs != nil {
|
||||
t.Fatalf("expected no errors, got: %v", errs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_partialLineAcrossWrites(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
// "[[[user]]]" split across two writes
|
||||
v.Validate([]byte("[["))
|
||||
v.Validate([]byte("[user]]]\nbody\n"))
|
||||
// Now validate the end marker separately
|
||||
v.Validate([]byte("[[[end]]]\n"))
|
||||
errs := v.Validate(nil) // flush remaining buffer
|
||||
if errs != nil {
|
||||
t.Fatalf("expected no errors, got: %v", errs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_fatalRejectsFurtherData(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
// First write causes a fatal nesting error
|
||||
v.Validate([]byte("[[[user]]]\n[[[assistant]]]\n"))
|
||||
errs := v.Validate([]byte("[[[end]]]\n"))
|
||||
if len(errs) == 0 {
|
||||
t.Fatal("expected rejection after fatal, got none")
|
||||
}
|
||||
if !strings.Contains(errs[0].Error(), "already failed") {
|
||||
t.Errorf("expected 'already failed' message, got: %s", errs[0].Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_hasViolations(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
if v.HasViolations() {
|
||||
t.Fatal("fresh validator should have no violations")
|
||||
}
|
||||
v.Validate([]byte("[[[user]]]\n"))
|
||||
if v.HasViolations() {
|
||||
t.Fatal("valid open block should not be a violation yet")
|
||||
}
|
||||
v.Validate([]byte("[[[assistant]]]\n"))
|
||||
if !v.HasViolations() {
|
||||
t.Fatal("nested block should trigger HasViolations")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseBlockHeader(t *testing.T) {
|
||||
tests := []struct {
|
||||
input string
|
||||
want *format.BlockHeader
|
||||
}{
|
||||
{"[[[user]]]", &format.BlockHeader{Role: "user"}},
|
||||
{"[[[assistant]]]", &format.BlockHeader{Role: "assistant"}},
|
||||
{"[[[tool:shell]]]", &format.BlockHeader{Role: "tool", Name: "shell"}},
|
||||
{"[[[call:file_read]]]", &format.BlockHeader{Role: "call", Name: "file_read"}},
|
||||
{"[[[retry]]]", &format.BlockHeader{Role: "retry"}},
|
||||
{"not a header", nil},
|
||||
{"[[[user]]] extra", nil},
|
||||
{"", nil},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
got := format.ParseBlockHeader(tt.input)
|
||||
if (got == nil) != (tt.want == nil) {
|
||||
t.Errorf("ParseBlockHeader(%q) = %v, want %v", tt.input, got, tt.want)
|
||||
continue
|
||||
}
|
||||
if got != nil && (got.Role != tt.want.Role || got.Name != tt.want.Name) {
|
||||
t.Errorf("ParseBlockHeader(%q) = [%s:%s], want [%s:%s]",
|
||||
tt.input, got.Role, got.Name, tt.want.Role, tt.want.Name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsEndMarker(t *testing.T) {
|
||||
tests := []struct {
|
||||
input string
|
||||
want bool
|
||||
}{
|
||||
{"[[[end]]]", true},
|
||||
{"[[[end]]] ", false},
|
||||
{" [[[end]]]", false},
|
||||
{"[[[end]]]\n", false},
|
||||
{"[[[user]]]", false},
|
||||
{"", false},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
got := format.IsEndMarker(tt.input)
|
||||
if got != tt.want {
|
||||
t.Errorf("IsEndMarker(%q) = %v, want %v", tt.input, got, tt.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidate_lineNumbers(t *testing.T) {
|
||||
v := NewLogValidator()
|
||||
data := []byte("[[[user]]]\nline2\n[[[assistant]]]\n")
|
||||
errs := v.Validate(data)
|
||||
if len(errs) == 0 {
|
||||
t.Fatal("expected nesting error at line 3, got none")
|
||||
}
|
||||
// The nested assistant header is on line 3
|
||||
if ve, ok := errs[0].(*violation); ok {
|
||||
if ve.line != 3 {
|
||||
t.Errorf("expected violation at line 3, got line %d", ve.line)
|
||||
}
|
||||
}
|
||||
}
|
||||
65
fs/types.go
65
fs/types.go
|
|
@ -2,12 +2,15 @@ package fs
|
|||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"ollie/agent"
|
||||
"ollie/backend"
|
||||
"ollie/format"
|
||||
"ollie/session"
|
||||
"ollie/toolsrv"
|
||||
)
|
||||
|
|
@ -139,10 +142,11 @@ type AgentLog struct {
|
|||
prevPrompt []byte
|
||||
remote string
|
||||
chatCond *sync.Cond
|
||||
validator *LogValidator // enforces [[[block]]]/[[[end]]] protocol structure
|
||||
}
|
||||
|
||||
func NewAgentLog(remote string) *AgentLog {
|
||||
al := &AgentLog{remote: remote}
|
||||
al := &AgentLog{remote: remote, validator: NewLogValidator()}
|
||||
al.chatCond = sync.NewCond(al.mu.RLocker())
|
||||
return al
|
||||
}
|
||||
|
|
@ -152,6 +156,10 @@ func (as *AgentLog) AppendLog(data []byte) {
|
|||
return
|
||||
}
|
||||
as.mu.Lock()
|
||||
errs := as.validator.Validate(data)
|
||||
for _, e := range errs {
|
||||
fmt.Fprintf(os.Stderr, "log structure violation: %v\n", e)
|
||||
}
|
||||
as.log = append(as.log, data...)
|
||||
as.logVers++
|
||||
as.mu.Unlock()
|
||||
|
|
@ -180,6 +188,25 @@ func (as *AgentLog) PrevPrompt() []byte { return as.prevPrompt }
|
|||
func (as *AgentLog) Mu() *sync.RWMutex { return &as.mu }
|
||||
func (as *AgentLog) Remote() string { return as.remote }
|
||||
|
||||
// Validator returns the log validator for external inspection.
|
||||
func (as *AgentLog) Validator() *LogValidator { return as.validator }
|
||||
|
||||
// UnclosedAtTeardown checks for unclosed blocks at agent teardown time.
|
||||
// Returns a list of violations if any blocks were opened but never closed,
|
||||
// or nil if all blocks are properly terminated.
|
||||
func (as *AgentLog) UnclosedAtTeardown() []error {
|
||||
as.mu.Lock()
|
||||
defer as.mu.Unlock()
|
||||
var errs []error
|
||||
for _, b := range as.validator.UnclosedBlocks() {
|
||||
errs = append(errs, fmt.Errorf(
|
||||
"unclosed block [%s:%s] opened at line %d — missing [[[end]]]",
|
||||
b.header.Role, b.header.Name, b.line,
|
||||
))
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
// NewEventHandler returns an event handler that writes to an AgentLog.
|
||||
func NewEventHandler(as *AgentLog) agent.EventHandler {
|
||||
streamingRole := ""
|
||||
|
|
@ -195,13 +222,13 @@ func NewEventHandler(as *AgentLog) agent.EventHandler {
|
|||
as.AppendLog([]byte("\n" + streamingFence + "\n"))
|
||||
streamingFence = ""
|
||||
}
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
}
|
||||
if ev.Role == "assistant" && ev.ResponseID != "" {
|
||||
as.AppendLog([]byte("[[[assistant:" + ev.ResponseID + "]]]\n"))
|
||||
as.AppendLog([]byte(format.AssistantHeaderPrefix + ev.ResponseID + "]]]\n"))
|
||||
streamingResponseID = ev.ResponseID
|
||||
} else {
|
||||
as.AppendLog([]byte("[[[" + ev.Role + "]]]\n"))
|
||||
as.AppendLog([]byte(format.RoleDelim(ev.Role) + "\n"))
|
||||
}
|
||||
if ev.Role == "assistant" {
|
||||
as.mu.Lock()
|
||||
|
|
@ -225,7 +252,7 @@ func NewEventHandler(as *AgentLog) agent.EventHandler {
|
|||
as.AppendLog([]byte("\n" + streamingFence + "\n"))
|
||||
streamingFence = ""
|
||||
}
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
streamingRole = ""
|
||||
} else {
|
||||
as.AppendLog([]byte(ev.Content))
|
||||
|
|
@ -237,12 +264,12 @@ func NewEventHandler(as *AgentLog) agent.EventHandler {
|
|||
as.AppendLog([]byte("\n" + streamingFence + "\n"))
|
||||
streamingFence = ""
|
||||
}
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
if ev.Content == "" {
|
||||
as.AppendLog([]byte("[[[tool:" + ev.Name + "]]]\n```" + ev.OutputFormat + "\n"))
|
||||
streamingFence = "```"
|
||||
as.AppendLog([]byte(format.ToolDelim(ev.Name) + ev.OutputFormat))
|
||||
streamingFence = format.ToolFence
|
||||
streamingRole = "tool"
|
||||
} else {
|
||||
as.AppendLog(session.FormatEvent(ev))
|
||||
|
|
@ -254,23 +281,23 @@ func NewEventHandler(as *AgentLog) agent.EventHandler {
|
|||
return
|
||||
}
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog([]byte("[[[retry]]]\n"))
|
||||
as.AppendLog([]byte(format.RetryTag + "\n"))
|
||||
streamingRole = "retry"
|
||||
as.AppendLog(session.FormatEvent(ev))
|
||||
|
||||
case "call":
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog(session.FormatEvent(ev))
|
||||
|
||||
default:
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte("\n" + format.EndTag + "\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog(session.FormatEvent(ev))
|
||||
|
|
@ -292,23 +319,23 @@ func replayMessagesToLog(as *AgentLog, messages []backend.Message) {
|
|||
switch m.Role {
|
||||
case "system":
|
||||
case "user":
|
||||
as.AppendLog([]byte("[[[user]]]\n" + m.Content + "\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte(format.UserTag + "\n" + m.Content + "\n" + format.EndTag + "\n"))
|
||||
case "assistant":
|
||||
header := "[[[assistant]]]"
|
||||
if m.ID != "" {
|
||||
header = "[[[assistant:" + m.ID + "]]]"
|
||||
as.AppendLog([]byte(format.AssistantHeaderPrefix + m.ID + "]]]\n"))
|
||||
} else {
|
||||
as.AppendLog([]byte(format.AssistantTag + "\n"))
|
||||
}
|
||||
as.AppendLog([]byte(header + "\n"))
|
||||
if m.Content != "" {
|
||||
as.AppendLog([]byte(m.Content + "\n"))
|
||||
}
|
||||
for _, tc := range m.ToolCalls {
|
||||
args := strings.Join(strings.Fields(string(tc.Arguments)), " ")
|
||||
as.AppendLog([]byte("[[[call:" + tc.Name + "]]]\n" + args + "\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte(format.CallDelim(tc.Name) + args + "\n" + format.EndTag + "\n"))
|
||||
}
|
||||
as.AppendLog([]byte("[[[end]]]\n"))
|
||||
as.AppendLog([]byte(format.EndTag + "\n"))
|
||||
case "tool":
|
||||
as.AppendLog([]byte("[[[tool]]]\n```\n" + strings.TrimRight(m.Content, "\n") + "\n```\n[[[end]]]\n"))
|
||||
as.AppendLog([]byte(format.ToolTag + "\n" + format.ToolFence + strings.TrimRight(m.Content, "\n") + "\n```\n" + format.EndTag + "\n"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
2
kde
2
kde
|
|
@ -1 +1 @@
|
|||
Subproject commit 08ec46027ee897cb34cd367590c50001a36667e0
|
||||
Subproject commit bb1454f85a4a72b43899cd5df078eb861533a86d
|
||||
Loading…
Reference in New Issue