| 1 | package shellrun |
| 2 | |
| 3 | import ( |
| 4 | "strings" |
| 5 | "unicode/utf8" |
| 6 | ) |
| 7 | |
| 8 | // Flush finishes a progress stream after the owner has drained the child, |
| 9 | // including cancellation. An unfinished rune is represented once, not silently |
| 10 | // lost or emitted as several invalid JSON strings. |
| 11 | func (w *progressWriter) Flush() { |
| 12 | w.mu.Lock() |
| 13 | defer w.mu.Unlock() |
| 14 | if w.emit == nil || w.truncated { |
| 15 | return |
| 16 | } |
| 17 | w.writeUTF8(nil, true) |
| 18 | } |
| 19 | |
| 20 | func (w *progressWriter) writeUTF8(p []byte, final bool) { |
| 21 | // At most one additional rune is needed to detect truncation. Do not copy |
| 22 | // an arbitrarily large Write when joining it to a pending character. |
| 23 | p = p[:min(len(p), max(0, w.limit-w.forwarded)+utf8.UTFMax)] |
| 24 | if len(w.pending) > 0 { |
| 25 | p = append(w.pending, p...) |
| 26 | w.pending = nil |
| 27 | } |
| 28 | var out strings.Builder |
| 29 | truncated := false |
| 30 | for len(p) > 0 { |
| 31 | if !utf8.FullRune(p) { |
| 32 | if !final { |
| 33 | w.pending = append([]byte(nil), p...) |
| 34 | break |
| 35 | } |
| 36 | p = []byte(string(utf8.RuneError)) |
| 37 | } |
| 38 | r, size := utf8.DecodeRune(p) |
| 39 | encodedSize := size |
| 40 | if r == utf8.RuneError && size == 1 { |
| 41 | encodedSize = 3 |
| 42 | } |
| 43 | if w.forwarded+out.Len()+encodedSize > w.limit { |
| 44 | truncated = true |
| 45 | break |
| 46 | } |
| 47 | out.WriteRune(r) |
| 48 | p = p[size:] |
| 49 | } |
| 50 | if out.Len() > 0 { |
| 51 | w.forwarded += out.Len() |
| 52 | w.emit(out.String()) |
| 53 | } |
| 54 | if truncated { |
| 55 | w.truncated = true |
| 56 | w.pending = nil |
| 57 | if w.marker != "" { |
| 58 | w.emit(w.marker) |
| 59 | } |
| 60 | } |
| 61 | } |
| 62 |