Pad/internal/editor/logic.go
Greg Pomerantz 83f7affee9 Fix word-wrap scroll jump: visual-line mapping via WrapIndex
Every scroll<->content mapping site (window start, sub-line shift, tap
mapping, max-scroll clamp, selection menu/handle positions) assumed
1 logical line = 1 visual line. When the viewport top crossed the
bottom of a wrapped line, the view jumped past the wrapped remainder
(jump magnitude (count-1)*lh) instead of moving pixel-by-pixel.

- WrapIndex (internal/editor/wrap_index.go): Fenwick tree of
  per-logical-line visual-line counts, parallel to the LineIndex;
  built at index-build time, bookkept by the same
  UpdateLineIndexAfter{Insert,Delete} hooks (never under-stale: every
  touched line resets to the estimate, the next shaping pass
  re-corrects it).
- scrollVisualDecompose: the scroll offset lives in visual-line space:
  k = LineForVisual(floor(s/lh)), r = s - V(k)*lh. All mapping sites
  go through it, so the viewport top is always exactly s into the
  document's visual space (V(k)*lh + r = s) — the jump invariant.
  All-ones index reduces to the legacy 1:1 mapping (pre-shaping and
  non-wrapped behavior unchanged by construction).
- Correction pipeline: the renderer's per-frame VisualLineStarts are
  grouped per logical line and written back (applyWrapCounts). The
  layout feedback now carries the exact window text the layout was
  shaped for (carried in the frame) plus the window start line and the
  content-edit counter; corrections apply only on edit-counter match,
  and grouping over the current window text (wrong after a scroll moved
  the window) is no longer possible.
- bytePosToScreenXY now applies the sub-line shift and the scaled line
  pitch: the selection menu/handles were off by up to a full line.
- maxScroll uses TotalVisuals() with the effective (font-scaled) line
  height; the bottom clamp lands exactly on the file end for wrapped
  content.
- VisibleByteRange returns the real start line (was hardcoded 0).
- emitFrame: replace the unread handoff frame with the newer snapshot
  instead of dropping it — a dropped final frame was never re-emitted
  (emission is event-driven), leaving the consumer one state behind
  forever; fixes the pre-existing TestRealFile_ShiftSelectionInsert
  failure. Still non-blocking.

Tests (mutation-verified where practical): wrap_index_test.go (Fenwick
vs naive model, 3000 ops), wrap_bookkeeping_test.go (edit hooks vs
shadow-string oracle, 400 ops — caught a real m=0 under-marking),
wrap_mapping_test.go (the jump regression: V(k)*lh + r == s over sweeps
+ random offsets; legacy-identity pin; boundary sweep), wrap_apply_test.go
(VisualLineStarts grouping + guards — the first version exposed the
always-true WindowStartByte guard that blocked all post-scroll
corrections). go test -race ./... green.

On-device (emulator, 60 wrapped lines): dp sweep 0/17/50/67/134/340
lands on LINE000-vl0/1/3, LINE001-vl0, LINE002-vl0, LINE005-vl0 —
pixel-exact 1:1, no jump (dp 134 is where the old code jumped to
LINE008); bottom clamp exact.

Docs: architecture.md §6.2 (visual-line space invariant),
development_plan.md (Phase 13), spec.md (wrap + clamp lines).
2026-08-17 19:46:47 -04:00

748 lines
25 KiB
Go

package editor
import (
"log"
"os"
"path/filepath"
"strconv"
"strings"
"sync"
"time"
"pad/internal/browser"
"pad/internal/io/pool"
"pad/internal/io/pool/mock"
"pad/internal/io/pool/types"
"pad/internal/ui"
)
// ConfigEvent represents a window configuration change (resize, orientation).
// PixelWidth and PixelHeight are the raw pixel dimensions from Gio.
type ConfigEvent struct {
PixelWidth int
PixelHeight int
}
// ScaleEvent represents a metric change (HiDPI scale factor).
type ScaleEvent struct {
Scale float32
// FontScale is the user font-size setting (Metric.PxPerSp / PxPerDp).
// The shaper draws baselines in sp, so the rendered line pitch in
// density-dp is EditorLineHeight()*FontScale; all geometry bookkeeping
// follows it (see EffectiveLineHeight). 0 = unknown, treated as 1.0.
FontScale float32
}
// ConfigUpdate is a common interface for all configuration updates.
// Both ConfigEvent and ScaleEvent implement this interface.
type ConfigUpdate interface {
apply(*State)
}
func (e ConfigEvent) apply(s *State) {
s.PixelWidth = e.PixelWidth
s.PixelHeight = e.PixelHeight
}
func (e ScaleEvent) apply(s *State) {
s.SetScale(e.Scale)
s.SetFontScale(e.FontScale)
}
// ResultEvent represents a completed async task result.
type ResultEvent struct {
// Future: add result fields here
}
// Logic runs the logic goroutine and provides channels for communication.
//
// Single-owner invariant (architecture.md §1): the logic goroutine is the
// sole reader/writer of l.state. Every other goroutine talks to it through
// the channels below. The only exception is Inspect, a test-only request
// channel whose fn still executes on the owner.
type Logic struct {
state *State
browserManager *browser.BrowserManager
configChan chan ConfigUpdate
frameChan chan Frame // frames carry the view-state snapshot
inputChan chan []ui.InputEvent
layoutChan chan ui.LayoutFeedback
resultChan chan ResultEvent
searchQueryChan chan string
openFileChan chan string
retryChan chan string // auto-save retries
autosaveChan chan struct{} // auto-save debounce ticks (timer -> owner)
inspectChan chan *inspectReq
// Per-file write protocol (see requestSave). Workers are a shared pool and
// the on-disk staging file is per-file, so two concurrent writes for the
// same file would interleave on the temp file; the protocol keeps at most
// one write in flight per file and re-issues deferred saves from the write
// result, so the rename that lands last always carries the newest content.
// All three maps are touched only on the owner goroutine.
writeInFlight map[string]int // filename -> fileVersion carried by the in-flight write
savePending map[string]bool // filename -> save requested while a write was in flight
retryScheduled map[string]time.Time // filename -> armed-but-unfired backoff retry
workerPool *pool.WorkerPool
mockFS pool.FileSystem
done chan struct{}
exitWg sync.WaitGroup
saveTimer *time.Timer // auto-save debounce timer; non-nil while pending
lastEmit time.Time // time of the last frame emission (profiler cadence)
debugCmdC chan string // one-shot debug commands from the cmd-file poller; nil = disabled
}
// NewLogic creates a new Logic instance, accepting an optional mockFS.
func NewLogic(mfs pool.FileSystem, path string, openfunc func(string)) *Logic {
state := NewState()
TheState = state
// The openfunc (e.g. the Android Termux bridge) is kept on TheState.open
// as an optional external hook, but it is NOT the tap path: tapping a file
// must open it in the in-app editor (doc/spec.md). The previous wiring set
// ui.OpenFile = openfunc, which routed every tap through the external
// bridge and bypassed the editor entirely.
TheState.open = openfunc
ui.OpenFile = func(path string) { OpenFile(path) }
// Initialize mock filesystem if nil
mockFS := mfs
if mockFS == nil {
mockFS = mock.NewFileSystem()
populateMockFileSystem(mockFS)
}
// Initialize worker pool
wp := pool.NewWorkerPool(8) // Increased worker pool size to 8
wp.Start()
state.Browser.CurrentPath = path
bm, _ := browser.NewBrowserManager(&state.Browser, wp, mockFS)
TheLogic = &Logic{
state: state,
browserManager: bm,
configChan: make(chan ConfigUpdate),
frameChan: make(chan Frame, 1),
inputChan: make(chan []ui.InputEvent),
layoutChan: make(chan ui.LayoutFeedback),
resultChan: make(chan ResultEvent),
searchQueryChan: make(chan string),
openFileChan: make(chan string),
retryChan: make(chan string, 1), // Buffered channel
autosaveChan: make(chan struct{}),
inspectChan: make(chan *inspectReq),
writeInFlight: make(map[string]int),
savePending: make(map[string]bool),
retryScheduled: make(map[string]time.Time),
workerPool: wp,
mockFS: mockFS,
done: make(chan struct{}),
}
return TheLogic
}
// ConfigChan returns the unified config channel for the logic goroutine.
// Accepts ConfigEvent (size) and ScaleEvent (scale factor).
func (l *Logic) ConfigChan() chan<- ConfigUpdate {
return l.configChan
}
// FrameChan returns the frame channel for the logic goroutine.
func (l *Logic) FrameChan() <-chan Frame {
return l.frameChan
}
// InputChan returns the input channel for the logic goroutine.
func (l *Logic) InputChan() chan<- []ui.InputEvent {
return l.inputChan
}
// LayoutChan returns the glyph layout feedback channel.
func (l *Logic) LayoutChan() chan<- ui.LayoutFeedback {
return l.layoutChan
}
// ResultChan returns the result channel for the logic goroutine.
func (l *Logic) ResultChan() chan<- ResultEvent {
return l.resultChan
}
// SearchQueryChan returns the search query channel for the logic goroutine.
// The main goroutine sends updated search text here when it detects a change.
func (l *Logic) SearchQueryChan() chan<- string {
return l.searchQueryChan
}
// ClipboardSetChan returns the channel the logic goroutine uses to request a
// system clipboard write; the main goroutine executes the Gio clipboard op.
func (l *Logic) ClipboardSetChan() <-chan string {
return l.state.clipboardSetChan
}
// PasteReqChan returns the channel the logic goroutine uses to request the
// system clipboard content for paste.
func (l *Logic) PasteReqChan() <-chan struct{} {
return l.state.pasteReqChan
}
// PasteChan delivers system clipboard content to the logic goroutine for
// paste (tests use it to inject clipboard text without a main loop).
func (l *Logic) PasteChan() chan<- string {
return l.state.pasteChan
}
var TheState *State
var TheLogic *Logic
// Run runs the logic goroutine loop.
func (l *Logic) Run() {
l.exitWg.Add(1)
defer l.exitWg.Done()
// Dispatch initial directory index build on startup
l.workerPool.Dispatch(pool.NewBuildIndexTask(l.state.Browser.CurrentPath, l.mockFS))
for {
select {
case <-l.done:
// Settle in-flight file writes before exiting: the post-exit
// synchronous FlushAll and workerPool.Stop must not race a
// straggling worker write, whose rename could otherwise land after
// the final flush and promote a stale snapshot.
l.drainWrites()
return
case update := <-l.configChan:
update.apply(l.state)
l.emitFrame()
case fb := <-l.layoutChan:
// Store the full GlyphLayout on editor state.
// Derive LastLineY from it for scroll clamping.
layout := fb.GlyphLayout
l.state.Editor.GlyphLayout = layout
var derivedLastLineY ui.Dp
if len(layout.Y) > 0 {
derivedLastLineY = layout.Y[len(layout.Y)-1]
}
// Wrap-count correction (see WrapIndex): the layout describes the
// window it was shaped for (fb.WindowStartLine); apply its visual
// line counts to those lines, but only if no edit has shifted the
// lines since shaping (fb.EditSeq correlates with the content).
if fb.EditSeq == l.state.Editor.EditSeq {
l.state.applyWrapCounts(fb)
}
if derivedLastLineY != l.state.LastLineY {
l.state.LastLineY = derivedLastLineY
l.emitFrame()
}
case events := <-l.inputChan:
for _, evt := range events {
evt.Handler(evt.Data)
}
l.emitFrame()
case query := <-l.searchQueryChan:
if query != l.state.Browser.Query {
l.state.Browser.Query = query
if l.state.page == BrowserPage {
browser.HandleSearch(&l.state.Browser, query)
}
}
l.emitFrame()
case path := <-l.openFileChan:
// Create chunked buffer for virtual scrolling
chunkSize := DefaultChunkSize
cb := NewChunkedBuffer(path, chunkSize, l.mockFS, "")
cb.SetWorkerPool(l.workerPool)
TheState.Editor.ChunkedBuffer = cb
// Dispatch stat task to get file size
l.workerPool.Dispatch(pool.NewStatFileTask(path, l.mockFS))
case filename := <-l.retryChan:
log.Printf("Logic: Retrying save for %s", filename)
delete(l.retryScheduled, filename)
l.requestSave(filename)
case <-l.autosaveChan:
// Auto-save debounce tick. The timer goroutine only sent a token;
// the owner snapshots content and dispatches the write.
l.saveTimer = nil
l.requestSave(l.state.Editor.Filename)
case p := <-l.state.pasteChan:
// Clipboard content arrived from the main goroutine: insert it
// (replacing any live selection, per the selection-aware edit rule).
HandlePaste(p)
l.emitFrame()
case req := <-l.inspectChan:
// Test-only: fn runs on the owner, preserving single ownership.
req.resp <- req.fn(l.state)
case res := <-l.workerPool.ResultChan():
l.handleWorkerResult(res)
case <-l.resultChan:
l.emitFrame()
case cmd := <-l.debugCmdC:
l.applyDebugCmd(cmd)
}
}
}
// emitFrame computes the current frame, records a profiler probe (if enabled),
// and hands it to the main goroutine. Centralizing emission here ensures the
// in-app profiler (PerfRecord) sees every frame exactly once, on the owner
// goroutine. Must be called on the logic goroutine.
func (l *Logic) emitFrame() {
elems := l.state.layout(l.browserManager)
now := time.Now()
if PerfRecord != nil {
var delta time.Duration
if !l.lastEmit.IsZero() {
delta = now.Sub(l.lastEmit)
}
l.lastEmit = now
s := l.state
rec := ProbeRecord{T: now, DeltaMs: float64(delta.Nanoseconds()) / 1e6, Page: pageName(s.page)}
if s.page == EditorPage {
rec.ScrollDP = float32(s.ScrollOffset)
rec.MaxScrollDP = float32(s.MaxScroll)
rec.VisStart = s.VisibleStart
rec.VisEnd = s.VisibleEnd
if cb := s.Editor.ChunkedBuffer; cb != nil {
if li := cb.LineIndex; li != nil {
rec.TotalLines = li.LineCount()
}
}
} else {
rec.ScrollDP = float32(s.Browser.ScrollOffset)
}
PerfRecord(rec)
}
f := l.frameOf(elems)
select {
case l.frameChan <- f:
default:
// Main has not consumed the previous frame yet. Frames are snapshots
// of the latest state, so the newest is always the most valuable:
// replace the unread one instead of dropping this one. A dropped
// final frame would never be re-emitted (emission is event-driven),
// leaving the consumer one state behind until the next event. Still
// non-blocking: a drain from a buffer whose consumer is gone (e.g.
// the post-done write drain) still just removes the unread frame.
select {
case <-l.frameChan:
default:
return // consumer gone; nothing to replace into
}
select {
case l.frameChan <- f:
default:
}
}
}
// EnableDebugCmdPoll starts a background poller (debug-only) that watches
// <dir>/cmd for a one-shot scroll command and forwards it to the owner for
// application. Used to jump the editor to specific scroll offsets for
// performance/clamping validation. The poller does the (blocking) file read
// off the owner and sends the command via debugCmdC; the owner applies it.
func (l *Logic) EnableDebugCmdPoll(dir string) {
l.debugCmdC = make(chan string, 1)
l.exitWg.Add(1)
go func() {
defer l.exitWg.Done()
t := time.NewTicker(120 * time.Millisecond)
defer t.Stop()
cmdPath := filepath.Join(dir, "cmd")
for {
select {
case <-l.done:
return
case <-t.C:
b, err := os.ReadFile(cmdPath)
if err != nil {
continue
}
cmd := strings.TrimSpace(string(b))
if cmd == "" {
continue
}
// Consume the command so it is applied exactly once.
_ = os.WriteFile(cmdPath, nil, 0o644)
select {
case l.debugCmdC <- cmd:
case <-l.done:
return
}
}
}
}()
}
// applyDebugCmd applies a one-shot debug scroll command from the cmd-file
// poller. Commands: "top", "bottom", "frac <0..1>", "dp <int>". Must be
// called on the logic goroutine.
func (l *Logic) applyDebugCmd(cmd string) {
s := l.state
if s.page != EditorPage {
log.Printf("DebugCmd: %q ignored (not on editor page)", cmd)
return
}
fields := strings.Fields(cmd)
if len(fields) == 0 {
return
}
var target ui.Dp
switch fields[0] {
case "top":
target = 0
case "bottom":
target = s.MaxScroll
case "frac":
if len(fields) < 2 {
return
}
f, err := strconv.ParseFloat(fields[1], 64)
if err != nil || f < 0 || f > 1 {
log.Printf("DebugCmd: bad frac %q", fields[1])
return
}
target = ui.Dp(float64(f) * float64(s.MaxScroll))
case "dp":
if len(fields) < 2 {
return
}
n, err := strconv.Atoi(fields[1])
if err != nil {
log.Printf("DebugCmd: bad dp %q", fields[1])
return
}
target = ui.Dp(n)
default:
log.Printf("DebugCmd: unknown %q", cmd)
return
}
if target < 0 {
target = 0
}
if target > s.MaxScroll {
target = s.MaxScroll
}
s.ScrollOffset = target
log.Printf("DebugCmd: %q -> scroll=%d maxScroll=%d", cmd, int(s.ScrollOffset), int(s.MaxScroll))
l.emitFrame()
}
// fullContentBytes reconstructs the full file content from the chunked
// buffer (or the deprecated full Buffer). Returns ok=false on error.
// Must be called on the logic goroutine.
func (l *Logic) fullContentBytes() ([]byte, bool) {
if l.state.Editor.ChunkedBuffer != nil {
fullContent, err := l.state.Editor.ChunkedBuffer.FullContent()
if err != nil {
log.Printf("Error reconstructing full content: %v", err)
return nil, false
}
return []byte(fullContent), true
}
return []byte(l.state.Editor.Buffer), true
}
// markDirty triggers the auto-save debounce timer.
// Must be called on the logic goroutine.
//
// The timer callback runs on a timer goroutine and only sends a token on
// autosaveChan; the owner (Run loop) does all state reads and the worker
// dispatch. This keeps the timer goroutine out of state (architecture.md §1).
func (l *Logic) markDirty() {
if l.state.Editor.Filename == "" {
return
}
l.state.Editor.fileVersion[l.state.Editor.Filename]++
// 1s debounce
if l.saveTimer != nil {
l.saveTimer.Stop()
}
l.saveTimer = time.AfterFunc(1*time.Second, func() {
l.autosaveChan <- struct{}{}
})
}
// handleWorkerResult processes results from the worker pool.
func (l *Logic) handleWorkerResult(res pool.Result) {
if res.IsBrowserResult() {
l.browserManager.HandleResult(res)
} else if res.TaskType == pool.TypeReadFile {
if res.Success {
if content, ok := res.Data.([]byte); ok {
if l.state.Editor.ChunkedBuffer != nil {
// Full-load the in-range file: set the whole content (split into
// resident chunks). No lazy load, so no stale-disk re-read.
l.state.Editor.ChunkedBuffer.SetFileSize(int64(len(content)))
l.state.Editor.ChunkedBuffer.SetContent(content)
// Also update the Buffer field for backward compatibility and for tests checking it.
l.state.Editor.Buffer = string(content)
} else {
// Fallback: populate the deprecated Buffer field
l.state.Editor.Buffer = string(content)
}
}
}
} else if res.TaskType == pool.TypeReadChunk {
// In-range files load fully via SetContent, so this is only a fallback.
if cb := l.state.Editor.ChunkedBuffer; cb != nil {
delete(cb.loadingChunks, res.ChunkIdx) // Use explicit ChunkIdx field
if res.Success {
if chunk, ok := res.Data.([]byte); ok {
for len(cb.chunks) <= res.ChunkIdx {
cb.chunks = append(cb.chunks, nil)
}
cb.chunks[res.ChunkIdx] = chunk
}
} else {
log.Printf("Logic: ReadChunkTask failed for %s chunk %d: %v", res.FilePath, res.ChunkIdx, res.Error)
}
}
} else if res.TaskType == pool.TypeStatFile {
if res.Success {
if stat, ok := res.Data.(*pool.FileStat); ok {
// Size guard: refuse to edit files above the limit. The browser can
// still list them; the editor shows a "too large to edit" notice.
if stat.Size > MaxEditableFileSize {
l.state.Editor.TooLarge = true
l.state.Editor.TooLargeSize = stat.Size
log.Printf("Logic: %s is %d bytes, exceeds the %d-byte edit limit", stat.Path, stat.Size, MaxEditableFileSize)
l.emitFrame()
return
}
if l.state.Editor.ChunkedBuffer != nil {
l.state.Editor.ChunkedBuffer.SetFileSize(stat.Size)
// Full-load: read the whole file and dispatch a line index build.
l.workerPool.Dispatch(pool.NewReadFileTask(stat.Path, l.mockFS))
l.workerPool.Dispatch(pool.NewBuildLineIndexTask(stat.Path, l.mockFS))
}
}
}
} else if res.TaskType == pool.TypeBuildLineIndex {
if res.Success {
if idx, ok := res.Data.(*types.LineIndex); ok {
if cb := l.state.Editor.ChunkedBuffer; cb != nil {
cb.LineIndex = idx
// The wrap index is created with the line index: same line
// set, same lifecycle. Counts start at the all-ones
// estimate and are corrected by shaped frames.
cb.WrapIndex = NewWrapIndex(idx.LineCount())
}
}
}
} else if res.TaskType == pool.TypeWriteFile {
f := res.FilePath
wrote, tracked := l.writeInFlight[f]
delete(l.writeInFlight, f)
pending := l.savePending[f]
delete(l.savePending, f)
if res.Success {
if tracked {
// Record the version whose content was actually written (the
// snapshot version), not the version now: edits made while the
// write was in flight sit ahead of it and trigger the re-issue
// below.
l.state.Editor.lastWriteVersion[f] = wrote
}
l.state.Editor.SetWriteFailed(f, false)
// Latest-state-wins re-issue: a save was deferred while this write
// was in flight, or an edit arrived during it — write the newer
// content now (only the active file is reconstructable).
if f == l.state.Editor.Filename && (pending || l.state.Editor.fileVersion[f] > l.state.Editor.lastWriteVersion[f]) {
l.requestSave(f)
}
} else {
log.Printf("Auto-save failed for %s: %v", f, res.Error)
l.state.Editor.SetWriteFailed(f, true)
// Trigger retry logic: exponential backoff
attempts := l.state.Editor.IncrementRetryAttempts(f)
// Simple backoff: 1s, 2s, 4s, 8s... max 30s
delay := time.Duration(1<<(attempts-1)) * time.Second
if delay > 30*time.Second {
delay = 30 * time.Second
}
l.retryScheduled[f] = time.Now().Add(delay)
filename := f
log.Printf("Scheduling retry for %s in %v (attempt %d)", filename, delay, attempts)
time.AfterFunc(delay, func() {
log.Printf("Firing retry for %s", filename)
select {
case l.retryChan <- filename:
default:
// A token for the same file is already queued; it will
// re-snapshot the latest content when it fires. Dropping
// keeps the timer goroutine from ever blocking.
log.Printf("Retry token for %s dropped: one already queued", filename)
}
})
}
}
l.emitFrame()
}
// applyBuildIndexResult applies a completed BuildIndexTask result to browser state.
func (l *Logic) applyBuildIndexResult(res pool.Result) {
// Delegated to browserManager
}
// applyLoadPagesResult applies a completed LoadPagesTask result to browser state.
func (l *Logic) applyLoadPagesResult(res pool.Result) {
// Delegated to browserManager
}
// applyReadDirResult applies a completed ReadDirTask result to browser state.
func (l *Logic) applyReadDirResult(res pool.Result) {
// Delegated to browserManager
}
// State returns the current state.
func (l *Logic) State() *State {
return l.state
}
// PruneMaps removes old entries to keep memory usage bounded.
// Keeps entries for the last 1000 files.
func (l *Logic) PruneMaps() {
if len(l.state.Editor.fileVersion) <= 1000 {
return
}
// Simply clear the maps for now. A true LRU would require
// tracking access times.
l.state.Editor.fileVersion = make(map[string]int)
l.state.Editor.lastWriteVersion = make(map[string]int)
l.state.Editor.writeFailed = make(map[string]bool)
l.state.Editor.retryAttempts = make(map[string]int)
}
// Done signals the logic goroutine to stop.
func (l *Logic) Done() {
close(l.done)
}
// WaitForExit blocks until the logic goroutine has fully stopped. After
// this returns, the caller (e.g. Shutdown) may touch state single-threaded.
func (l *Logic) WaitForExit() {
l.exitWg.Wait()
}
// FlushAll triggers synchronous writes for all dirty files.
//
// Must be called either from the logic goroutine (e.g. via GoToBrowser) or
// after the logic goroutine has fully stopped (Shutdown). It touches state
// directly, so it must never run concurrently with Run (single-owner
// invariant, architecture.md §1).
func (l *Logic) FlushAll() {
for filename := range l.state.Editor.fileVersion {
if l.state.Editor.IsDirty() && l.state.Editor.Filename == filename {
if _, inflight := l.writeInFlight[filename]; inflight {
// A worker is already writing this file and we cannot block the
// owner (we are it): defer. The write's result handler re-issues
// a fresh write while the file is still dirty, so the disk ends
// up with the latest content.
l.savePending[filename] = true
continue
}
content, ok := l.fullContentBytes()
if !ok {
continue
}
// In a real app, this would be a blocking call to the FS
l.mockFS.WriteFileAtomic(filename, content)
l.state.Editor.lastWriteVersion[filename] = l.state.Editor.fileVersion[filename]
}
}
}
// requestSave is the single entry point for "persist the active file now"
// (autosave tick, retry token, deferred re-issue). It implements the
// per-file write protocol:
//
// - at most one write per file is in flight at any time. The worker pool is
// shared and the on-disk staging file is per-file, so two concurrent
// writes for the same file would interleave on the staging file and could
// rename a byte-mixture into place;
// - each write carries the file version at the moment its content was
// snapshotted (writeInFlight[f]); on success exactly that version is
// recorded as written, so any edit that arrived while the write was in
// flight leaves the file dirty and the result handler re-issues a fresh
// write with the newer content;
// - a save requested while a write is in flight is deferred (savePending)
// and re-issued by the result handler, so "last rename wins" always
// coincides with "newest snapshot wins".
//
// Must be called on the owner goroutine.
func (l *Logic) requestSave(filename string) {
if filename == "" || filename != l.state.Editor.Filename {
// Only the active file's content lives in memory (its ChunkedBuffer);
// a save for any previously opened file cannot be reconstructed.
return
}
if _, inflight := l.writeInFlight[filename]; inflight {
l.savePending[filename] = true
return
}
content, ok := l.fullContentBytes()
if !ok {
return
}
l.writeInFlight[filename] = l.state.Editor.fileVersion[filename]
l.workerPool.DispatchNonBlocking(pool.NewWriteFileTask(filename, content, l.mockFS))
}
const (
// writeDrainTimeout bounds the post-done wait for in-flight writes so a
// pathologically failing write can never hang shutdown.
writeDrainTimeout = 5 * time.Second
// writeDrainTick is the idle poll interval of the drain loop.
writeDrainTick = 20 * time.Millisecond
)
// drainWrites waits for in-flight file writes (and armed retries) to settle
// after the done signal, so that Shutdown's post-exit synchronous FlushAll
// and workerPool.Stop cannot race a straggling worker write. The write
// protocol may re-issue follow-up writes while draining; the deadline bounds
// the wait. Must be called on the owner goroutine.
func (l *Logic) drainWrites() {
if len(l.writeInFlight) == 0 && len(l.retryScheduled) == 0 {
return
}
log.Printf("Logic: draining %d in-flight write(s) before exit", len(l.writeInFlight))
deadline := time.Now().Add(writeDrainTimeout)
for {
if len(l.writeInFlight) == 0 && len(l.retryScheduled) == 0 {
return
}
if time.Now().After(deadline) {
log.Printf("Logic: write drain timed out with %d in flight; exiting (per-write unique temp files keep every rename a complete snapshot)", len(l.writeInFlight))
return
}
select {
case res := <-l.workerPool.ResultChan():
l.handleWorkerResult(res)
case filename := <-l.retryChan:
delete(l.retryScheduled, filename)
l.requestSave(filename)
case <-l.autosaveChan:
l.saveTimer = nil
l.requestSave(l.state.Editor.Filename)
case <-time.After(writeDrainTick):
// No events: re-check settle and deadline.
}
}
}
// Shutdown gracefully shuts down the logic goroutine and worker pool.
// The logic goroutine is stopped FIRST so that FlushAll can touch state
// single-threaded; the worker pool is stopped last so in-flight results
// still have a reader until then.
func (l *Logic) Shutdown() {
l.Done()
l.WaitForExit()
l.FlushAll()
l.workerPool.Stop()
}