mirror of
https://github.com/jesseduffield/lazygit.git
synced 2026-09-10 07:36:27 -04:00
ReadToEnd reads the rest of a view's content on the render task's own goroutine, and calls back once it has. Nothing held a task for that, so lazygit counted as idle from the moment the caller returned until the callback ran. The search prompt in the focused main view opens from such a callback, so an integration test takes the idle report as its cue to carry on, and presses its next key while the prompt is not open yet. Hold the task in ReadToEnd rather than in the caller, so that every caller is covered (see docs/dev/Busy.md). Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
710 lines
24 KiB
Go
710 lines
24 KiB
Go
package tasks
|
|
|
|
import (
|
|
"bufio"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"os/exec"
|
|
"strconv"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/jesseduffield/lazygit/pkg/commands/oscommands"
|
|
"github.com/jesseduffield/lazygit/pkg/gocui"
|
|
"github.com/jesseduffield/lazygit/pkg/utils"
|
|
"github.com/sasha-s/go-deadlock"
|
|
"github.com/sirupsen/logrus"
|
|
)
|
|
|
|
// Cmd abstracts over a started external process. *exec.Cmd satisfies the bulk
|
|
// of it via ExecCmd, but pty implementations can supply their own types — on
|
|
// Windows, ConPTY has to spawn via CreateProcess directly and can't use
|
|
// *exec.Cmd (see golang/go#62708).
|
|
type Cmd interface {
|
|
Wait() error
|
|
String() string
|
|
// Terminate makes the process stop early, as gracefully as the platform
|
|
// allows. It doesn't wait for the process to exit.
|
|
Terminate() error
|
|
}
|
|
|
|
// ExecCmd adapts *exec.Cmd to Cmd.
|
|
type ExecCmd struct {
|
|
*exec.Cmd
|
|
}
|
|
|
|
// Terminate sends SIGTERM on Unix. On Windows it does nothing, so a stopped
|
|
// command keeps running until it next writes to its (by then closed) output
|
|
// pipe.
|
|
func (c ExecCmd) Terminate() error {
|
|
return oscommands.TerminateProcessGracefully(c.Process)
|
|
}
|
|
|
|
// This file revolves around running commands that will be output to the main panel
|
|
// in the gui. If we're flicking through the commits panel, we want to invoke a
|
|
// `git show` command for each commit, but we don't want to read the entire output
|
|
// at once (because that would slow things down); we just want to fill the panel
|
|
// and then read more as the user scrolls down. We also want to ensure that we're only
|
|
// ever running one `git show` command at time, and that we only have one command
|
|
// writing its output to the main panel at a time.
|
|
|
|
const THROTTLE_TIME = time.Millisecond * 30
|
|
|
|
// we use this to check if the system is under stress right now. Hopefully this makes sense on other machines
|
|
const COMMAND_START_THRESHOLD = time.Millisecond * 10
|
|
|
|
type ViewBufferManager struct {
|
|
// this blocks until the task has been properly stopped
|
|
stopCurrentTask func()
|
|
|
|
// this is what we write the output of the task to. It's typically a view
|
|
writer io.Writer
|
|
|
|
waitingMutex deadlock.Mutex
|
|
// Guards newTaskID and taskKey, which identify the most recently requested
|
|
// task. Both are written on the goroutine NewTask spawns, and taskKey is
|
|
// read from the UI thread (GetTaskKey), so neither may be touched without
|
|
// holding this.
|
|
taskIDMutex deadlock.Mutex
|
|
Log *logrus.Entry
|
|
newTaskID int
|
|
// The requests by which the currently-running task is told to read more lines
|
|
// (e.g. as the user scrolls), and which it answers once it has. The task
|
|
// serving them comes and goes; see readRequestQueue.
|
|
readRequests *readRequestQueue
|
|
taskKey string
|
|
|
|
// Resets the view's scroll position to the top. A render whose content is
|
|
// different from what the view last showed (a different command key) calls
|
|
// this — but at its *first paint*, not when the task starts: the off-screen
|
|
// render leaves the previous content displayed until the swap, so resetting
|
|
// the origin up front would scroll that still-displayed content to the top
|
|
// before the new content replaces it. See newContentPending.
|
|
resetOrigin func()
|
|
|
|
// Whether the content the running task is rendering differs from what the
|
|
// view is currently showing (i.e. the command key changed). Two things key
|
|
// off it: the loading indicator only takes the view over when it is set,
|
|
// since there is no point clearing content we are about to render
|
|
// identically; and the first paint that reveals the content resets the
|
|
// scroll to the top and clears it.
|
|
//
|
|
// It deliberately outlives the task that set it: a task can be stopped and
|
|
// replaced before it ever paints — a background refresh landing just after
|
|
// the user clicked a different item, say — and the replacement, which
|
|
// renders the same content and so sets nothing of its own, still has to do
|
|
// what that task was owed.
|
|
newContentPending atomic.Bool
|
|
|
|
// Whether a command task is currently reading content into the view. While
|
|
// this is true the content is still growing, so callers (e.g. the layout)
|
|
// must not clamp the view's scroll position to the amount loaded so far.
|
|
loading atomic.Bool
|
|
|
|
// beforeStart is the function that is called before starting a new task
|
|
beforeStart func()
|
|
refreshView func()
|
|
onEndOfInput func()
|
|
|
|
// beginRender starts an off-screen render: the new content is built without
|
|
// disturbing what's displayed. swapInRender then promotes it to the display
|
|
// in one step. Together they keep the view showing the previous render until
|
|
// the new one has read enough to paint, instead of revealing it line by line.
|
|
beginRender func()
|
|
swapInRender func()
|
|
|
|
// see docs/dev/Busy.md
|
|
// A gocui task is not the same thing as the tasks defined in this file.
|
|
// A gocui task simply represents the fact that lazygit is busy doing something,
|
|
// whereas the tasks in this file are about rendering content to a view.
|
|
newGocuiTask func() gocui.Task
|
|
|
|
// Runs f on the UI thread and blocks until it has completed. All mutations
|
|
// of the view happen through this, so that the view is only ever touched on
|
|
// the UI thread (where it is also laid out and drawn), never on the task's
|
|
// own goroutine.
|
|
onUIThread func(f func()) error
|
|
|
|
// if the user flicks through a heap of items, with each one
|
|
// spawning a process to render something to the main view,
|
|
// it can slow things down quite a bit. In these situations we
|
|
// want to throttle the spawning of processes. Atomic because it's set
|
|
// from one task's stop goroutine and read when the next task starts.
|
|
throttle atomic.Bool
|
|
}
|
|
|
|
type LinesToRead struct {
|
|
// The total number of lines the task should have read once this request is
|
|
// satisfied. This is an absolute count from the start of the task, not a
|
|
// delta: the task keeps track of how many lines it has already read and only
|
|
// reads the shortfall, so a request for a total at or below what has already
|
|
// been read reads nothing. -1 means read all the way to the end.
|
|
Total int
|
|
|
|
// Number of lines after which we have read enough to fill the view, and can
|
|
// do an initial refresh. Only set for the initial read request; -1 for
|
|
// subsequent requests.
|
|
InitialRefreshAfter int
|
|
|
|
// Function to call after reading the lines is done
|
|
Then func()
|
|
}
|
|
|
|
func (self *ViewBufferManager) GetTaskKey() string {
|
|
self.taskIDMutex.Lock()
|
|
defer self.taskIDMutex.Unlock()
|
|
|
|
return self.taskKey
|
|
}
|
|
|
|
func NewViewBufferManager(
|
|
log *logrus.Entry,
|
|
writer io.Writer,
|
|
beforeStart func(),
|
|
refreshView func(),
|
|
onEndOfInput func(),
|
|
resetOrigin func(),
|
|
beginRender func(),
|
|
swapInRender func(),
|
|
newGocuiTask func() gocui.Task,
|
|
onUIThread func(f func()) error,
|
|
) *ViewBufferManager {
|
|
return &ViewBufferManager{
|
|
readRequests: newReadRequestQueue(),
|
|
Log: log,
|
|
writer: writer,
|
|
beforeStart: beforeStart,
|
|
refreshView: refreshView,
|
|
onEndOfInput: onEndOfInput,
|
|
resetOrigin: resetOrigin,
|
|
beginRender: beginRender,
|
|
swapInRender: swapInRender,
|
|
newGocuiTask: newGocuiTask,
|
|
onUIThread: onUIThread,
|
|
}
|
|
}
|
|
|
|
// ReadLines asks the task to ensure it has read at least totalLines lines in
|
|
// total. Because the count is absolute rather than a delta, repeated requests
|
|
// (e.g. as the user scrolls down, back up, and down again) don't re-read lines
|
|
// that have already been read: the task only ever reads the shortfall.
|
|
func (self *ViewBufferManager) ReadLines(totalLines int) {
|
|
// A request with no Then needs no answer, so there is nothing to do when no
|
|
// task is there to take it.
|
|
self.readRequests.enqueue(LinesToRead{Total: totalLines, InitialRefreshAfter: -1})
|
|
}
|
|
|
|
// IsLoading reports whether a command task is currently reading content into the
|
|
// view, meaning the content is still growing.
|
|
func (self *ViewBufferManager) IsLoading() bool {
|
|
return self.loading.Load()
|
|
}
|
|
|
|
// StartLoading marks the view as loading content. It must be called
|
|
// synchronously when a command/pty task is started, before the task's goroutine
|
|
// runs, so that a layout pass happening in between doesn't clamp the scroll
|
|
// position to the not-yet-loaded content. It is cleared when the task reaches
|
|
// the end of its input.
|
|
func (self *ViewBufferManager) StartLoading() {
|
|
self.loading.Store(true)
|
|
}
|
|
|
|
func (self *ViewBufferManager) ReadToEnd(then func()) {
|
|
// The reading happens on the task's own goroutine, and the caller hears about
|
|
// it through then, so lazygit must not count as idle in between.
|
|
task := self.newGocuiTask()
|
|
answered := func() {
|
|
task.Done()
|
|
if then != nil {
|
|
then()
|
|
}
|
|
}
|
|
|
|
request := LinesToRead{Total: -1, InitialRefreshAfter: -1, Then: answered}
|
|
if !self.readRequests.enqueue(request) {
|
|
// With no task reading, everything there is to read has been read.
|
|
answered()
|
|
}
|
|
}
|
|
|
|
// stopServingReadRequests takes the task away from the read-request queue and
|
|
// answers whatever it never got to, so that nobody is left waiting for a callback
|
|
// that isn't coming.
|
|
func (self *ViewBufferManager) stopServingReadRequests() {
|
|
for _, request := range self.readRequests.stopServing() {
|
|
if request.Then != nil {
|
|
request.Then()
|
|
}
|
|
}
|
|
}
|
|
|
|
func (self *ViewBufferManager) NewCmdTask(start func() (Cmd, io.Reader), prefix string, linesToRead LinesToRead, onDoneFn func()) func(TaskOpts) error {
|
|
return func(opts TaskOpts) error {
|
|
var onDoneOnce sync.Once
|
|
var onFirstPageShownOnce sync.Once
|
|
|
|
onFirstPageShown := func() {
|
|
onFirstPageShownOnce.Do(func() {
|
|
opts.InitialContentLoaded()
|
|
})
|
|
}
|
|
|
|
onDone := func() {
|
|
if onDoneFn != nil {
|
|
onDoneOnce.Do(onDoneFn)
|
|
}
|
|
onFirstPageShown()
|
|
}
|
|
|
|
if self.throttle.Load() {
|
|
self.Log.Info("throttling task")
|
|
time.Sleep(THROTTLE_TIME)
|
|
}
|
|
|
|
select {
|
|
case <-opts.Stop:
|
|
onDone()
|
|
return nil
|
|
default:
|
|
}
|
|
|
|
startTime := time.Now()
|
|
cmd, r := start()
|
|
timeToStart := time.Since(startTime)
|
|
|
|
done := make(chan struct{})
|
|
|
|
go utils.Safe(func() {
|
|
select {
|
|
case <-done:
|
|
// The command finished and did not have to be preemptively stopped before the next command.
|
|
// No need to throttle.
|
|
self.throttle.Store(false)
|
|
case <-opts.Stop:
|
|
// we use the time it took to start the program as a way of checking if things
|
|
// are running slow at the moment. This is admittedly a crude estimate, but
|
|
// the point is that we only want to throttle when things are running slow
|
|
// and the user is flicking through a bunch of items.
|
|
self.throttle.Store(time.Since(startTime) < THROTTLE_TIME && timeToStart > COMMAND_START_THRESHOLD)
|
|
|
|
// Kill the still-running command. The only reason to do this is to save CPU usage
|
|
// when flicking through several very long diffs when diff.algorithm = histogram is
|
|
// being used, in which case multiple git processes continue to calculate expensive
|
|
// diffs in the background even though they have been stopped already.
|
|
if err := cmd.Terminate(); err != nil {
|
|
self.Log.Errorf("error when trying to terminate cmd task: %v; Command: %v", err, cmd.String())
|
|
}
|
|
|
|
// close the task's stdout pipe (or the pty if we're using one) to make the command terminate
|
|
onDone()
|
|
}
|
|
})
|
|
|
|
loadingMutex := deadlock.Mutex{}
|
|
|
|
// Begin serving before any goroutine starts, so that the first request below
|
|
// can't arrive before there is a task to take it.
|
|
readRequests := self.readRequests.beginServing()
|
|
|
|
scanner := bufio.NewScanner(r)
|
|
scanner.Split(utils.ScanLinesAndTruncateWhenLongerThanBuffer(bufio.MaxScanTokenSize))
|
|
|
|
lineChan := make(chan []byte)
|
|
lineWrittenChan := make(chan struct{})
|
|
|
|
// We're reading from the scanner in a separate goroutine because on windows
|
|
// if running git through a shim, we sometimes kill the parent process without
|
|
// killing its children, meaning the scanner blocks forever. This solution
|
|
// leaves us with a dead goroutine, but it's better than blocking all
|
|
// rendering to main views.
|
|
go utils.Safe(func() {
|
|
defer close(lineChan)
|
|
for scanner.Scan() {
|
|
select {
|
|
case <-opts.Stop:
|
|
return
|
|
case lineChan <- scanner.Bytes():
|
|
// We need to confirm the data has been fed into the view before we
|
|
// pull more from the scanner because the scanner uses the same backing
|
|
// array and we don't want to be mutating that while it's being written
|
|
<-lineWrittenChan
|
|
}
|
|
}
|
|
|
|
if err := scanner.Err(); err != nil {
|
|
self.Log.Error(err)
|
|
}
|
|
})
|
|
|
|
loaded := false
|
|
|
|
go utils.Safe(func() {
|
|
ticker := time.NewTicker(time.Millisecond * 200)
|
|
defer ticker.Stop()
|
|
select {
|
|
case <-opts.Stop:
|
|
return
|
|
case <-ticker.C:
|
|
loadingMutex.Lock()
|
|
// Only take the view over to say "loading..." when the content coming
|
|
// is different from what's on screen. A re-render of the same content
|
|
// leaves the view showing exactly what it should already, so clearing
|
|
// it for the message and then rendering the same thing back is a
|
|
// visible flicker for nothing — and a slow re-render of unchanged
|
|
// content is common (a background refresh over a repo with submodules
|
|
// that have uncommitted changes, say). The pending flag isn't consumed
|
|
// here; the first paint still owes the scroll reset.
|
|
if !loaded && self.newContentPending.Load() {
|
|
self.beforeStart()
|
|
// beforeStart cleared the previous content to show "loading...", so
|
|
// put the view back at the top for it (beforeStart doesn't touch the
|
|
// origin). The origin is view state the UI thread reads while laying
|
|
// out, so write it there.
|
|
_ = self.onUIThread(self.resetOrigin)
|
|
_, _ = self.writer.Write([]byte("loading..."))
|
|
self.refreshView()
|
|
}
|
|
loadingMutex.Unlock()
|
|
}
|
|
})
|
|
|
|
go utils.Safe(func() {
|
|
isViewStale := true
|
|
writeToView := func(content []byte) {
|
|
isViewStale = true
|
|
_, _ = self.writer.Write(content)
|
|
}
|
|
refreshViewIfStale := func() {
|
|
if isViewStale {
|
|
self.refreshView()
|
|
isViewStale = false
|
|
}
|
|
}
|
|
|
|
// Go's select picks randomly among ready cases, so once opts.Stop is
|
|
// closed the selects below could still service a ready data channel
|
|
// instead of bailing. Check stop explicitly first to give it priority:
|
|
// a task that's been stopped (it's being replaced by a newer one) must
|
|
// not touch the view here — it would start an off-screen render and
|
|
// write the prefix into it, clobbering what the incoming task is about
|
|
// to render.
|
|
stopped := func() bool {
|
|
select {
|
|
case <-opts.Stop:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// The total number of lines we have read so far. Requests specify an
|
|
// absolute target total (see LinesToRead.Total), so we compare against
|
|
// this to work out how many more lines, if any, we still need to read.
|
|
linesRead := 0
|
|
|
|
// The first paint swaps the off-screen render in to reveal the new
|
|
// content, and settles the scroll position in the same step — so the new
|
|
// content first appears already where it belongs, and no draw can land
|
|
// between the two and show it at the previous render's scroll. It happens
|
|
// once, either when we've read far enough (below) or at end of input for
|
|
// content shorter than that. Callers run it on the UI thread: it writes
|
|
// the view's origin.
|
|
painted := false
|
|
firstPaint := func() {
|
|
if painted {
|
|
return
|
|
}
|
|
painted = true
|
|
self.swapInRender()
|
|
if self.newContentPending.Swap(false) {
|
|
self.resetOrigin()
|
|
}
|
|
}
|
|
|
|
// Set LAZYGIT_SLOW_RENDER=<milliseconds> to sleep that long after each
|
|
// line is written to the view, stretching async loads out so the frames
|
|
// of a re-render become visible. Useful for debugging scroll/flicker
|
|
// behaviour; has no effect when the variable is unset.
|
|
var slowRenderPerLine time.Duration
|
|
if v := os.Getenv("LAZYGIT_SLOW_RENDER"); v != "" {
|
|
if ms, err := strconv.Atoi(v); err == nil {
|
|
slowRenderPerLine = time.Duration(ms) * time.Millisecond
|
|
}
|
|
}
|
|
|
|
outer:
|
|
for {
|
|
if stopped() {
|
|
break outer
|
|
}
|
|
linesToRead, ok := self.readRequests.dequeue()
|
|
if !ok {
|
|
// Nothing to read yet: wait to be told there is, or to be stopped.
|
|
select {
|
|
case <-opts.Stop:
|
|
break outer
|
|
case <-readRequests:
|
|
}
|
|
continue
|
|
}
|
|
{
|
|
callThen := func() {
|
|
if linesToRead.Then != nil {
|
|
linesToRead.Then()
|
|
}
|
|
}
|
|
for linesToRead.Total == -1 || linesRead < linesToRead.Total {
|
|
if stopped() {
|
|
callThen()
|
|
break outer
|
|
}
|
|
var ok bool
|
|
var line []byte
|
|
select {
|
|
case <-opts.Stop:
|
|
callThen()
|
|
break outer
|
|
case line, ok = <-lineChan:
|
|
// process line below
|
|
}
|
|
|
|
loadingMutex.Lock()
|
|
if !loaded {
|
|
// Build the new content off-screen, leaving the previous render
|
|
// displayed until we swap in below; this is what keeps an async
|
|
// re-render from showing a half-loaded buffer.
|
|
self.beginRender()
|
|
if prefix != "" {
|
|
writeToView([]byte(prefix))
|
|
}
|
|
loaded = true
|
|
}
|
|
loadingMutex.Unlock()
|
|
|
|
if !ok {
|
|
// lineChan is closed. At a genuine end of input we swap in what we
|
|
// read and finalize. But lineChan is also closed when this task has
|
|
// been stopped to make way for a newer one: stopping closes
|
|
// opts.Stop, and the scanner goroutine then closes lineChan, so the
|
|
// select above can land here instead of on the opts.Stop case. A
|
|
// stopped task is being replaced and must leave the view to the
|
|
// incoming task — swapping in its half-read buffer, clamping the
|
|
// origin, or clearing `loading` would all corrupt what that task is
|
|
// about to render. So bail out here, the same as the explicit stop
|
|
// case above.
|
|
select {
|
|
case <-opts.Stop:
|
|
callThen()
|
|
break outer
|
|
default:
|
|
}
|
|
// Genuine end of input: do the first paint now if it hasn't happened
|
|
// yet (the content was shorter than a screenful, so we never reached
|
|
// the point below), and flush the stale content. onEndOfInput reads
|
|
// the view's dimensions (to decide whether to scroll) and sets the
|
|
// origin, both of which are UI-thread-only, so run it there — as is
|
|
// firstPaint, which also writes the origin.
|
|
_ = self.onUIThread(func() {
|
|
firstPaint()
|
|
self.onEndOfInput()
|
|
})
|
|
// The content is fully loaded now, so it's safe again for the
|
|
// layout to clamp the scroll position to it. We deliberately
|
|
// don't clear this when stopped (rather than EOF'd), because that
|
|
// means a newer task is taking over and is still loading.
|
|
self.loading.Store(false)
|
|
callThen()
|
|
break outer
|
|
}
|
|
writeToView(append(line, '\n'))
|
|
lineWrittenChan <- struct{}{}
|
|
linesRead++
|
|
|
|
if slowRenderPerLine > 0 {
|
|
time.Sleep(slowRenderPerLine)
|
|
}
|
|
|
|
if linesRead == linesToRead.InitialRefreshAfter {
|
|
// We have read enough lines to fill the view, so do the first paint
|
|
// and refresh to show it. Continue reading and refresh again at the
|
|
// end to make sure the scrollbar has the right size.
|
|
_ = self.onUIThread(firstPaint)
|
|
refreshViewIfStale()
|
|
}
|
|
}
|
|
refreshViewIfStale()
|
|
onFirstPageShown()
|
|
callThen()
|
|
}
|
|
}
|
|
|
|
// Whoever made a request the loop never got to is waiting to hear that the
|
|
// content it asked for has been read, and there is nothing here to read it
|
|
// any more: at end of input it has all been read already, and a task that
|
|
// was stopped is handing the view over to the one replacing it.
|
|
self.stopServingReadRequests()
|
|
|
|
refreshViewIfStale()
|
|
|
|
select {
|
|
case <-opts.Stop:
|
|
// If we stopped the task, don't block waiting for it; this could cause a delay if
|
|
// the process takes a while until it actually terminates. We still want to call
|
|
// Wait to reclaim any resources, but do it on a background goroutine, and ignore
|
|
// any errors.
|
|
go func() { _ = cmd.Wait() }()
|
|
default:
|
|
if err := cmd.Wait(); err != nil {
|
|
self.Log.Errorf("Unexpected error when running cmd task: %v; Failed command: %v", err, cmd.String())
|
|
}
|
|
}
|
|
|
|
// calling this here again in case the program ended on its own accord
|
|
onDone()
|
|
|
|
close(done)
|
|
close(lineWrittenChan)
|
|
})
|
|
|
|
self.readRequests.enqueue(linesToRead)
|
|
|
|
<-done
|
|
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// Close closes the task manager, killing whatever task may currently be running
|
|
func (self *ViewBufferManager) Close() {
|
|
// stopCurrentTask is written by NewTask's goroutine under waitingMutex (and
|
|
// so is the sync.Once it closes over), so read it under the lock and call
|
|
// the captured value; a task starting on shutdown must not race us here.
|
|
self.waitingMutex.Lock()
|
|
stopCurrentTask := self.stopCurrentTask
|
|
self.waitingMutex.Unlock()
|
|
|
|
if stopCurrentTask == nil {
|
|
return
|
|
}
|
|
|
|
c := make(chan struct{})
|
|
|
|
go utils.Safe(func() {
|
|
stopCurrentTask()
|
|
c <- struct{}{}
|
|
})
|
|
|
|
select {
|
|
case <-c:
|
|
return
|
|
case <-time.After(3 * time.Second):
|
|
fmt.Println("cannot kill child process")
|
|
}
|
|
}
|
|
|
|
// different kinds of tasks:
|
|
// 1) command based, where the manager can be asked to read more lines, but the command can be killed
|
|
// 2) string based, where the manager can also be asked to read more lines
|
|
|
|
type TaskOpts struct {
|
|
// Channel that tells the task to stop, because another task wants to run.
|
|
Stop chan struct{}
|
|
|
|
// Only for tasks which are long-running, where we read more lines sporadically.
|
|
// We use this to keep track of when a user's action is complete (i.e. all views
|
|
// have been refreshed to display the results of their action)
|
|
InitialContentLoaded func()
|
|
}
|
|
|
|
func (self *ViewBufferManager) NewTask(f func(TaskOpts) error, key string) error {
|
|
gocuiTask := self.newGocuiTask()
|
|
|
|
var completeTaskOnce sync.Once
|
|
|
|
completeGocuiTask := func() {
|
|
completeTaskOnce.Do(func() {
|
|
gocuiTask.Done()
|
|
})
|
|
}
|
|
|
|
// Assign the taskID synchronously so it reflects NewTask call order
|
|
// rather than the order in which the spawned goroutines happen to be
|
|
// scheduled. Otherwise two NewTask calls in quick succession can have
|
|
// their goroutines race, with the later-called task ending up with the
|
|
// lower taskID and losing the staleness check below.
|
|
self.taskIDMutex.Lock()
|
|
self.newTaskID++
|
|
taskID := self.newTaskID
|
|
self.taskIDMutex.Unlock()
|
|
|
|
go utils.Safe(func() {
|
|
defer completeGocuiTask()
|
|
|
|
self.taskIDMutex.Lock()
|
|
|
|
// Bail out before touching shared view state if a newer task has
|
|
// already been queued: if we reset the view here we'd do it for a task
|
|
// that's about to exit, potentially wiping output the winning task has
|
|
// already written.
|
|
if taskID < self.newTaskID {
|
|
self.taskIDMutex.Unlock()
|
|
return
|
|
}
|
|
|
|
// Note we don't reset the origin here even when the command key changed:
|
|
// that's deferred to the first paint that reveals the new content (see
|
|
// newContentPending), so the previous content — left displayed until the
|
|
// swap — doesn't visibly jump to the top before the new content appears.
|
|
// Read taskKey directly: we already hold the mutex that guards it, and
|
|
// GetTaskKey would take it again.
|
|
if self.taskKey != key && self.resetOrigin != nil {
|
|
self.newContentPending.Store(true)
|
|
}
|
|
self.taskKey = key
|
|
|
|
self.taskIDMutex.Unlock()
|
|
|
|
self.waitingMutex.Lock()
|
|
|
|
// Re-check staleness after acquiring waitingMutex: a newer task
|
|
// may have arrived while we were blocked here.
|
|
self.taskIDMutex.Lock()
|
|
if taskID < self.newTaskID {
|
|
self.waitingMutex.Unlock()
|
|
self.taskIDMutex.Unlock()
|
|
return
|
|
}
|
|
self.taskIDMutex.Unlock()
|
|
|
|
if self.stopCurrentTask != nil {
|
|
self.stopCurrentTask()
|
|
}
|
|
|
|
// Nothing serves read requests between one task and the next.
|
|
self.stopServingReadRequests()
|
|
|
|
stop := make(chan struct{})
|
|
notifyStopped := make(chan struct{})
|
|
|
|
var once sync.Once
|
|
onStop := func() {
|
|
close(stop)
|
|
<-notifyStopped
|
|
}
|
|
|
|
self.stopCurrentTask = func() { once.Do(onStop) }
|
|
|
|
self.waitingMutex.Unlock()
|
|
|
|
if err := f(TaskOpts{Stop: stop, InitialContentLoaded: completeGocuiTask}); err != nil {
|
|
self.Log.Error(err) // might need an onError callback
|
|
}
|
|
|
|
close(notifyStopped)
|
|
})
|
|
|
|
return nil
|
|
}
|