Fix stream chunk truncation; style TUI initial state from state file
server.go: the stream read handler truncated each chunk to the client's requested Count and advanced waitBase, silently dropping bytes beyond the buffer (e.g. a chat block larger than 64KB). It now caches each stream chunk per-fid (readCache + streamBase) and pages within it by absolute offset, fetching the next chunk only once the current one is fully drained, preserving 9P offset semantics. o tui: read the state file for the initial status bar and apply the same color styling and status text the event loop uses, so the bar is correct immediately on attach; the event stream drives subsequent updates.
This commit is contained in:
parent
89fbb7438a
commit
cd26595cb7
|
|
@ -65,15 +65,16 @@ func (g *GroupTable) IsMember(group, user string) bool {
|
||||||
}
|
}
|
||||||
|
|
||||||
type fid struct {
|
type fid struct {
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
path string
|
path string
|
||||||
qid plan9.Qid
|
qid plan9.Qid
|
||||||
mode uint8
|
mode uint8
|
||||||
entry fs.File
|
entry fs.File
|
||||||
writeBuf []byte
|
writeBuf []byte
|
||||||
readCache []byte // cached content for offset-based paging
|
readCache []byte // cached content for offset-based paging
|
||||||
waitBase string
|
waitBase string
|
||||||
dirCache []byte
|
streamBase int64 // absolute byte offset where readCache begins (stream mode)
|
||||||
|
dirCache []byte
|
||||||
|
|
||||||
// Event stream subscription (for event file with filter)
|
// Event stream subscription (for event file with filter)
|
||||||
eventFilter string
|
eventFilter string
|
||||||
|
|
@ -684,31 +685,53 @@ func read(root *fs.Tree, cs *connState, fc *plan9.Fcall, ctx context.Context) *p
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Default stream handling (no filter)
|
// Default stream handling (no filter).
|
||||||
|
//
|
||||||
|
// A stream chunk may exceed the client's requested Count. Cache the
|
||||||
|
// chunk per-fid and page within it by absolute offset; only fetch the
|
||||||
|
// next chunk once the current one is fully drained. This preserves 9P
|
||||||
|
// offset semantics and never drops bytes beyond Count.
|
||||||
cs.mu.RLock()
|
cs.mu.RLock()
|
||||||
f, fidOK := cs.fids[fc.Fid]
|
f, fidOK := cs.fids[fc.Fid]
|
||||||
var base string
|
|
||||||
if fidOK {
|
|
||||||
f.mu.Lock()
|
|
||||||
base = f.waitBase
|
|
||||||
f.mu.Unlock()
|
|
||||||
}
|
|
||||||
cs.mu.RUnlock()
|
cs.mu.RUnlock()
|
||||||
|
if !fidOK {
|
||||||
|
return errFcall(fc, "bad fid")
|
||||||
|
}
|
||||||
|
|
||||||
|
f.mu.Lock()
|
||||||
|
cached := f.readCache
|
||||||
|
streamBase := f.streamBase
|
||||||
|
f.mu.Unlock()
|
||||||
|
|
||||||
|
// Serve from the current cached chunk if the requested offset falls
|
||||||
|
// within it.
|
||||||
|
if cached != nil && int64(fc.Offset) >= streamBase && int64(fc.Offset) < streamBase+int64(len(cached)) {
|
||||||
|
local := int(int64(fc.Offset) - streamBase)
|
||||||
|
end := local + int(fc.Count)
|
||||||
|
if end > len(cached) {
|
||||||
|
end = len(cached)
|
||||||
|
}
|
||||||
|
data := cached[local:end]
|
||||||
|
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Current chunk drained (or first read): fetch the next chunk.
|
||||||
|
f.mu.Lock()
|
||||||
|
base := f.waitBase
|
||||||
|
f.mu.Unlock()
|
||||||
waitCtx, waitCancel := context.WithCancel(ctx)
|
waitCtx, waitCancel := context.WithCancel(ctx)
|
||||||
defer waitCancel()
|
defer waitCancel()
|
||||||
content, nextBase, err := entry.BlockingRead(waitCtx, base)
|
content, nextBase, err := entry.BlockingRead(waitCtx, base)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errFcall(fc, err.Error())
|
return errFcall(fc, err.Error())
|
||||||
}
|
}
|
||||||
|
f.mu.Lock()
|
||||||
if nextBase != "" {
|
if nextBase != "" {
|
||||||
cs.mu.Lock()
|
f.waitBase = nextBase
|
||||||
if f, ok := cs.fids[fc.Fid]; ok {
|
|
||||||
f.mu.Lock()
|
|
||||||
f.waitBase = nextBase
|
|
||||||
f.mu.Unlock()
|
|
||||||
}
|
|
||||||
cs.mu.Unlock()
|
|
||||||
}
|
}
|
||||||
|
f.readCache = content
|
||||||
|
f.streamBase = int64(fc.Offset)
|
||||||
|
f.mu.Unlock()
|
||||||
count := int(fc.Count)
|
count := int(fc.Count)
|
||||||
if count > len(content) {
|
if count > len(content) {
|
||||||
count = len(content)
|
count = len(content)
|
||||||
|
|
|
||||||
|
|
@ -524,10 +524,20 @@ cmd_tui() {
|
||||||
sleep 0.2
|
sleep 0.2
|
||||||
tmux send-keys -t "$tmux_session" "o '$ctx' prompt" Enter
|
tmux send-keys -t "$tmux_session" "o '$ctx' prompt" Enter
|
||||||
|
|
||||||
# Get initial state
|
# Get initial state directly from the state file, then style it the same
|
||||||
local initial_state
|
# way the event loop below does, so the status bar is correct immediately.
|
||||||
|
local initial_state initial_status initial_bg
|
||||||
initial_state=$(ollie-9p read "session/$s/agent/$a/state" 2>/dev/null || echo "unknown")
|
initial_state=$(ollie-9p read "session/$s/agent/$a/state" 2>/dev/null || echo "unknown")
|
||||||
tmux set-option -t "$tmux_session" status-right "$initial_state | %H:%M"
|
case "$initial_state" in
|
||||||
|
idle) initial_bg="colour70" ;;
|
||||||
|
calling*) initial_bg="colour214" ;;
|
||||||
|
thinking) initial_bg="colour33" ;;
|
||||||
|
paused) initial_bg="colour245" ;;
|
||||||
|
*) initial_bg="default" ;;
|
||||||
|
esac
|
||||||
|
initial_status=$(ollie-9p read "${PREFIX}/status" 2>/dev/null | head -n1)
|
||||||
|
tmux set-option -t "$tmux_session" status-right-style "fg=default,bg=$initial_bg"
|
||||||
|
tmux set-option -t "$tmux_session" status-right "${initial_status:-$initial_state} | %H:%M"
|
||||||
tmux set-option -t "$tmux_session" status-right-length 64
|
tmux set-option -t "$tmux_session" status-right-length 64
|
||||||
|
|
||||||
# Get session and agent IDs for event subscription
|
# Get session and agent IDs for event subscription
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue