Compare commits

...
Author SHA1 Message Date
pashpashpash 8e1ae0ebc3 markdown streaming for cli wip 2025-10-07 02:13:37 -07:00
9 changed files with 1231 additions and 13 deletions
+23
View File
@@ -14,10 +14,33 @@ require (
replace github.com/cline/grpc-go => ../src/generated/grpc-go
require (
github.com/alecthomas/chroma/v2 v2.14.0 // indirect
github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect
github.com/aymerick/douceur v0.2.0 // indirect
github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc // indirect
github.com/charmbracelet/glamour v0.10.0 // indirect
github.com/charmbracelet/lipgloss v1.1.1-0.20250404203927-76690c660834 // indirect
github.com/charmbracelet/x/ansi v0.8.0 // indirect
github.com/charmbracelet/x/cellbuf v0.0.13 // indirect
github.com/charmbracelet/x/exp/slice v0.0.0-20250327172914-2fdc97757edf // indirect
github.com/charmbracelet/x/term v0.2.1 // indirect
github.com/dlclark/regexp2 v1.11.0 // indirect
github.com/gorilla/css v1.0.1 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/lucasb-eyer/go-colorful v1.2.0 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
github.com/mattn/go-runewidth v0.0.16 // indirect
github.com/microcosm-cc/bluemonday v1.0.27 // indirect
github.com/muesli/reflow v0.3.0 // indirect
github.com/muesli/termenv v0.16.0 // indirect
github.com/rivo/uniseg v0.4.7 // indirect
github.com/spf13/pflag v1.0.5 // indirect
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
github.com/yuin/goldmark v1.7.8 // indirect
github.com/yuin/goldmark-emoji v1.0.5 // indirect
golang.org/x/net v0.41.0 // indirect
golang.org/x/sys v0.33.0 // indirect
golang.org/x/term v0.32.0 // indirect
golang.org/x/text v0.26.0 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20250707201910-8d1bb00bc6a7 // indirect
)
+51
View File
@@ -1,6 +1,28 @@
github.com/alecthomas/chroma/v2 v2.14.0 h1:R3+wzpnUArGcQz7fCETQBzO5n9IMNi13iIs46aU4V9E=
github.com/alecthomas/chroma/v2 v2.14.0/go.mod h1:QolEbTfmUHIMVpBqxeDnNBj2uoeI4EbYP4i6n68SG4I=
github.com/atotto/clipboard v0.1.4 h1:EH0zSVneZPSuFR11BlR9YppQTVDbh5+16AmcJi4g1z4=
github.com/atotto/clipboard v0.1.4/go.mod h1:ZY9tmq7sm5xIbd9bOK4onWV4S6X0u6GY7Vn0Yu86PYI=
github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k=
github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8=
github.com/aymerick/douceur v0.2.0 h1:Mv+mAeH1Q+n9Fr+oyamOlAkUNPWPlA8PPGR0QAaYuPk=
github.com/aymerick/douceur v0.2.0/go.mod h1:wlT5vV2O3h55X9m7iVYN0TBM0NH/MmbLnd30/FjWUq4=
github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc h1:4pZI35227imm7yK2bGPcfpFEmuY1gc2YSTShr4iJBfs=
github.com/charmbracelet/colorprofile v0.2.3-0.20250311203215-f60798e515dc/go.mod h1:X4/0JoqgTIPSFcRA/P6INZzIuyqdFY5rm8tb41s9okk=
github.com/charmbracelet/glamour v0.10.0 h1:MtZvfwsYCx8jEPFJm3rIBFIMZUfUJ765oX8V6kXldcY=
github.com/charmbracelet/glamour v0.10.0/go.mod h1:f+uf+I/ChNmqo087elLnVdCiVgjSKWuXa/l6NU2ndYk=
github.com/charmbracelet/lipgloss v1.1.1-0.20250404203927-76690c660834 h1:ZR7e0ro+SZZiIZD7msJyA+NjkCNNavuiPBLgerbOziE=
github.com/charmbracelet/lipgloss v1.1.1-0.20250404203927-76690c660834/go.mod h1:aKC/t2arECF6rNOnaKaVU6y4t4ZeHQzqfxedE/VkVhA=
github.com/charmbracelet/x/ansi v0.8.0 h1:9GTq3xq9caJW8ZrBTe0LIe2fvfLR/bYXKTx2llXn7xE=
github.com/charmbracelet/x/ansi v0.8.0/go.mod h1:wdYl/ONOLHLIVmQaxbIYEC/cRKOQyjTkowiI4blgS9Q=
github.com/charmbracelet/x/cellbuf v0.0.13 h1:/KBBKHuVRbq1lYx5BzEHBAFBP8VcQzJejZ/IA3iR28k=
github.com/charmbracelet/x/cellbuf v0.0.13/go.mod h1:xe0nKWGd3eJgtqZRaN9RjMtK7xUYchjzPr7q6kcvCCs=
github.com/charmbracelet/x/exp/slice v0.0.0-20250327172914-2fdc97757edf h1:rLG0Yb6MQSDKdB52aGX55JT1oi0P0Kuaj7wi1bLUpnI=
github.com/charmbracelet/x/exp/slice v0.0.0-20250327172914-2fdc97757edf/go.mod h1:B3UgsnsBZS/eX42BlaNiJkD1pPOUa+oF1IYC6Yd2CEU=
github.com/charmbracelet/x/term v0.2.1 h1:AQeHeLZ1OqSXhrAWpYUtZyX1T3zVxfpZuEQMIQaGIAQ=
github.com/charmbracelet/x/term v0.2.1/go.mod h1:oQ4enTYFV7QN4m0i9mzHrViD7TQKvNEEkHUMCmsxdUg=
github.com/cpuguy83/go-md2man/v2 v2.0.3/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o=
github.com/dlclark/regexp2 v1.11.0 h1:G/nrcoOa7ZXlpoa/91N3X7mM3r8eIlMBBJZvsz/mxKI=
github.com/dlclark/regexp2 v1.11.0/go.mod h1:DHkYz0B9wPfa6wondMfaivmHpzrQ3v9q8cnmRbL6yW8=
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
@@ -11,15 +33,41 @@ github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gorilla/css v1.0.1 h1:ntNaBIghp6JmvWnxbZKANoLyuXTPZ4cAMlo6RyhlbO8=
github.com/gorilla/css v1.0.1/go.mod h1:BvnYkspnSzMmwRK+b8/xgNPLiIuNZr6vbZBTPQ2A3b0=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/lucasb-eyer/go-colorful v1.2.0 h1:1nnpGOrhyZZuNyfu1QjKiUICQ74+3FNCN69Aj6K7nkY=
github.com/lucasb-eyer/go-colorful v1.2.0/go.mod h1:R4dSotOR9KMtayYi1e77YzuveK+i7ruzyGqttikkLy0=
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
github.com/mattn/go-runewidth v0.0.12/go.mod h1:RAqKPSqVFrSLVXbA8x7dzmKdmGzieGRCM46jaSJTDAk=
github.com/mattn/go-runewidth v0.0.16 h1:E5ScNMtiwvlvB5paMFdw9p4kSQzbXFikJ5SQO6TULQc=
github.com/mattn/go-runewidth v0.0.16/go.mod h1:Jdepj2loyihRzMpdS35Xk/zdY8IAYHsh153qUoGf23w=
github.com/mattn/go-sqlite3 v1.14.24 h1:tpSp2G2KyMnnQu99ngJ47EIkWVmliIizyZBfPrBWDRM=
github.com/mattn/go-sqlite3 v1.14.24/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y=
github.com/microcosm-cc/bluemonday v1.0.27 h1:MpEUotklkwCSLeH+Qdx1VJgNqLlpY2KXwXFM08ygZfk=
github.com/microcosm-cc/bluemonday v1.0.27/go.mod h1:jFi9vgW+H7c3V0lb6nR74Ib/DIB5OBs92Dimizgw2cA=
github.com/muesli/reflow v0.3.0 h1:IFsN6K9NfGtjeggFP+68I4chLZV2yIKsXJFNZ+eWh6s=
github.com/muesli/reflow v0.3.0/go.mod h1:pbwTDkVPibjO2kyvBQRBxTWEEGDGq0FlB1BIKtnHY/8=
github.com/muesli/termenv v0.16.0 h1:S5AlUN9dENB57rsbnkPyfdGuWIlkmzJjbFf0Tf5FWUc=
github.com/muesli/termenv v0.16.0/go.mod h1:ZRfOIKPFDYQoDFF4Olj7/QJbW60Ol/kL1pU3VfY/Cnk=
github.com/rivo/uniseg v0.1.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc=
github.com/rivo/uniseg v0.2.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc=
github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ=
github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88=
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
github.com/spf13/cobra v1.8.0 h1:7aJaZx1B85qltLMc546zn58BxxfZdR/W22ej9CFoEf0=
github.com/spf13/cobra v1.8.0/go.mod h1:WXLWApfZ71AjXPya3WOlMsY9yMs7YeiHhFVlvLyhcho=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
github.com/spf13/pflag v1.0.5/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM=
github.com/yuin/goldmark v1.7.1/go.mod h1:uzxRWxtg69N339t3louHJ7+O03ezfj6PlliRlaOzY1E=
github.com/yuin/goldmark v1.7.8 h1:iERMLn0/QJeHFhxSt3p6PeN9mGnvIKSpG9YYorDMnic=
github.com/yuin/goldmark v1.7.8/go.mod h1:uzxRWxtg69N339t3louHJ7+O03ezfj6PlliRlaOzY1E=
github.com/yuin/goldmark-emoji v1.0.5 h1:EMVWyCGPlXJfUXBXpuMu+ii3TIaxbVBnEX9uaDC4cIk=
github.com/yuin/goldmark-emoji v1.0.5/go.mod h1:tTkZEbwu5wkPmgTcitqddVxY9osFZiavD+r4AzQrh1U=
go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=
go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A=
go.opentelemetry.io/otel v1.37.0 h1:9zhNfelUvx0KBfu/gb+ZgeAfAgtWrfHJZcAqFC228wQ=
@@ -34,8 +82,11 @@ go.opentelemetry.io/otel/trace v1.37.0 h1:HLdcFNbRQBE2imdSEgm/kwqmQj1Or1l/7bW6mx
go.opentelemetry.io/otel/trace v1.37.0/go.mod h1:TlgrlQ+PtQO5XFerSPUYG0JSgGyryXewPGyayAWSBS0=
golang.org/x/net v0.41.0 h1:vBTly1HeNPEn3wtREYfy4GZ/NECgw2Cnl+nK6Nz3uvw=
golang.org/x/net v0.41.0/go.mod h1:B/K4NNqkfmg07DQYrbwvSluqCJOOXwUjeb/5lOisjbA=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw=
golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/term v0.32.0 h1:DR4lr0TjUs3epypdhTOkMmuF5CDFJ/8pOnbzMZPQ7bg=
golang.org/x/term v0.32.0/go.mod h1:uZG1FhGx848Sqfsq4/DlJr3xGGsYMu/L5GW4abiaEPQ=
golang.org/x/text v0.26.0 h1:P42AVeLghgTYr4+xUnTRKDMqpar+PtX7KWuNQL21L8M=
golang.org/x/text v0.26.0/go.mod h1:QK15LZJUUQVJxhz7wXgxSy/CJaTFjd0G+YLonydOVQA=
gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk=
+39
View File
@@ -5,6 +5,7 @@ import (
"strings"
"time"
"github.com/charmbracelet/glamour"
"github.com/cline/cli/pkg/cli/global"
"github.com/cline/cli/pkg/cli/types"
"github.com/cline/grpc-go/cline"
@@ -35,6 +36,44 @@ func (r *Renderer) RenderMessage(timestamp, prefix, text string) error {
return nil
}
// RenderTextWithMarkdown renders text with markdown styling in rich mode, plain text otherwise
func (r *Renderer) RenderTextWithMarkdown(timestamp, text string) error {
if text == "" {
return nil
}
cleanText := r.sanitizeText(text)
if cleanText == "" {
return nil
}
// Check if we're in rich mode
if global.Config.OutputFormat == "rich" {
// Render markdown using Glamour
renderer, err := glamour.NewTermRenderer(
glamour.WithStylePath("dark"),
glamour.WithWordWrap(80),
)
if err != nil {
// Fallback to plain rendering if Glamour fails
return r.RenderMessage(timestamp, "🤖", cleanText)
}
rendered, err := renderer.Render(cleanText)
if err != nil {
// Fallback to plain rendering if rendering fails
return r.RenderMessage(timestamp, "🤖", cleanText)
}
// Print timestamp and then the rendered markdown
fmt.Printf("[%s] 🤖:\n%s\n", timestamp, rendered)
return nil
}
// Fall back to plain text rendering
return r.RenderMessage(timestamp, "ASST TEXT", cleanText)
}
// RenderCommand renders a command execution
func (r *Renderer) RenderCommand(timestamp, command string, isExecuting bool) error {
if isExecuting {
+72 -6
View File
@@ -6,24 +6,36 @@ import (
"strings"
"sync"
"github.com/cline/cli/pkg/cli/global"
"github.com/cline/cli/pkg/cli/types"
"github.com/cline/cli/pkg/markdown"
)
// StreamingDisplay manages streaming message display with deduplication
type StreamingDisplay struct {
mu sync.RWMutex
state *types.ConversationState
renderer *Renderer
dedupe *MessageDeduplicator
mu sync.RWMutex
state *types.ConversationState
renderer *Renderer
dedupe *MessageDeduplicator
markdownRenderer *markdown.StreamRenderer
}
// NewStreamingDisplay creates a new streaming display manager
func NewStreamingDisplay(state *types.ConversationState, renderer *Renderer) *StreamingDisplay {
return &StreamingDisplay{
sd := &StreamingDisplay{
state: state,
renderer: renderer,
dedupe: NewMessageDeduplicator(),
}
// Initialize markdown renderer if in rich mode
if global.Config != nil && global.Config.OutputFormat == "rich" {
if mdRenderer, err := markdown.NewStreamRenderer(); err == nil {
sd.markdownRenderer = mdRenderer
}
}
return sd
}
// HandlePartialMessage processes partial messages with streaming support
@@ -83,7 +95,7 @@ func (s *StreamingDisplay) handleStreamingAsk(msg *types.ClineMessage, messageKe
// handleStreamingSay handles streaming SAY messages
func (s *StreamingDisplay) handleStreamingSay(msg *types.ClineMessage, messageKey, timestamp string, streamingMsg *types.StreamingMessage) error {
switch msg.Say {
case string(types.SayTypeText), string(types.SayTypeCompletionResult):
case string(types.SayTypeText), string(types.SayTypeCompletionResult), string(types.SayTypeReasoning):
return s.handleStreamingText(msg, messageKey, timestamp, streamingMsg)
case string(types.SayTypeCommand):
return s.handleStreamingCommand(msg, messageKey, timestamp, streamingMsg)
@@ -111,6 +123,57 @@ func (s *StreamingDisplay) handleStreamingText(msg *types.ClineMessage, messageK
return nil // Duplicate - ignore it
}
// Rich mode: use markdown streaming
if s.markdownRenderer != nil {
return s.handleMarkdownText(msg, messageKey, cleanText, streamingMsg)
}
// Plain mode: use typewriter
return s.handlePlainText(msg, messageKey, timestamp, cleanText, streamingMsg)
}
// handleMarkdownText handles text messages in rich mode with markdown rendering
func (s *StreamingDisplay) handleMarkdownText(msg *types.ClineMessage, messageKey, cleanText string, streamingMsg *types.StreamingMessage) error {
// Check if this is an update to the same message
if streamingMsg.CurrentKey == messageKey {
// Incremental update
if len(cleanText) > len(streamingMsg.LastText) && strings.HasPrefix(cleanText, streamingMsg.LastText) {
newChars := cleanText[len(streamingMsg.LastText):]
if err := s.markdownRenderer.WriteIncremental(newChars); err != nil {
return fmt.Errorf("failed to write incremental markdown: %w", err)
}
s.state.SetStreamingMessage(messageKey, cleanText)
} else {
// Non-incremental change - this shouldn't happen often in streaming
// Flush and start new
s.finishCurrentStream()
if err := s.markdownRenderer.WriteIncremental(cleanText); err != nil {
return fmt.Errorf("failed to write markdown: %w", err)
}
s.state.SetStreamingMessage(messageKey, cleanText)
}
} else {
// New message - flush previous and start new
s.finishCurrentStream()
if err := s.markdownRenderer.WriteIncremental(cleanText); err != nil {
return fmt.Errorf("failed to write markdown: %w", err)
}
s.state.SetStreamingMessage(messageKey, cleanText)
}
// If message is complete, flush the markdown
if !msg.Partial {
if err := s.markdownRenderer.FlushMessage(); err != nil {
return fmt.Errorf("failed to flush markdown: %w", err)
}
s.state.SetStreamingMessage("", "")
}
return nil
}
// handlePlainText handles text messages in plain mode with typewriter
func (s *StreamingDisplay) handlePlainText(msg *types.ClineMessage, messageKey, timestamp, cleanText string, streamingMsg *types.StreamingMessage) error {
// Check if this is an update to the same message
if streamingMsg.CurrentKey == messageKey {
// Show incremental changes
@@ -446,4 +509,7 @@ func (s *StreamingDisplay) Cleanup() {
if s.dedupe != nil {
s.dedupe.Stop()
}
if s.markdownRenderer != nil {
s.markdownRenderer.Close()
}
}
+8 -6
View File
@@ -146,13 +146,13 @@ func (h *SayHandler) handleText(msg *types.ClineMessage, dc *DisplayContext, tim
return nil
}
// Special case for the user's task input
prefix := "ASST TEXT"
// Special case for the user's task input (no markdown rendering)
if dc.MessageIndex == 0 {
prefix = "USER"
return dc.Renderer.RenderMessage(timestamp, "USER", msg.Text)
}
return dc.Renderer.RenderMessage(timestamp, prefix, msg.Text)
// For rich mode, render markdown. For other modes, use plain text
return dc.Renderer.RenderTextWithMarkdown(timestamp, msg.Text)
}
// handleReasoning handles reasoning messages
@@ -161,7 +161,8 @@ func (h *SayHandler) handleReasoning(msg *types.ClineMessage, dc *DisplayContext
return nil
}
return dc.Renderer.RenderMessage(timestamp, "THINKING", msg.Text)
// Render reasoning/thinking messages with markdown in rich mode
return dc.Renderer.RenderTextWithMarkdown(timestamp, msg.Text)
}
func (h *SayHandler) handleCompletionResult(msg *types.ClineMessage, dc *DisplayContext, timestamp string) error {
@@ -171,7 +172,8 @@ func (h *SayHandler) handleCompletionResult(msg *types.ClineMessage, dc *Display
text = strings.TrimSuffix(text, "HAS_CHANGES")
}
return dc.Renderer.RenderMessage(timestamp, "RESULT", text)
// Render completion results with markdown in rich mode
return dc.Renderer.RenderTextWithMarkdown(timestamp, text)
}
// handleUserFeedback handles user feedback messages
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2024 Charmbracelet, Inc.
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+823
View File
@@ -0,0 +1,823 @@
// Package flow provides memory-efficient streaming markdown rendering with configurable flow control.
//
// The flow package enables processing large markdown files without loading them entirely into memory.
// It supports three flow modes for different use cases:
//
// - Unbuffered (-1): Aggressive flushing at safe boundaries for minimal latency
// - Buffered (0): Maximum buffering until EOF for optimal throughput
// - Windowed (N): Bounded buffering with N-byte window for balanced performance
//
// Key features:
// - Smart boundary detection to avoid breaking markdown structures
// - Code fence awareness to prevent splitting within code blocks
// - YAML frontmatter detection and proper handling
// - Buffer overflow protection with configurable limits
// - Production-ready error handling and resource management
//
// Example usage:
//
// render := func(data []byte) ([]byte, error) {
// // Your markdown renderer (e.g., glamour)
// return glamour.Render(string(data))
// }
//
// err := Flow(ctx, reader, writer, 0, render) // Buffered mode
// if err != nil {
// log.Fatal(err)
// }
package flow
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"syscall"
)
const (
// Core modes
Unbuffered = -1 // Aggressive flushing at boundaries
Buffered = 0 // No flushing until EOF
Windowed = 1048576 // 1 MiB max buffer length
)
// RenderFunc is a function that renders markdown to output bytes
type RenderFunc func(data []byte) ([]byte, error)
// Config configures the streaming buffer behavior
type Config struct {
Window int64 // -1: minimal buffering (aggressive flush), 0: maximum buffering (no flush until EOF), N: bounded at N bytes
Render RenderFunc // glamour renderer (byte-based)
}
// Validate checks config parameters for safety
func (c Config) Validate() error {
if c.Render == nil {
return errors.New("render function cannot be nil")
}
if c.Window > Windowed {
return fmt.Errorf("window size %d exceeds maximum %d", c.Window, Windowed)
}
if c.Window < -1 {
return fmt.Errorf("invalid window size %d (must be >= -1)", c.Window)
}
return nil
}
// Buffer handles incremental markdown rendering with configurable buffering
type Buffer struct {
config Config
writer io.Writer // output writer
accum []byte // accumulated markdown
written int64 // total bytes written
hadSplits bool // track if we actually split at buffer boundaries
pendingBlank []byte // deferred glamour blank lines awaiting context
offset int // -1 unknown, 0 none, >0 end index of detected frontmatter
// Fence tracking for smart boundaries
inFence bool // currently inside ``` code fence
lastFence int // Track position where fence closed
postFence bool // Just exited a fence block
}
// NewBuffer creates a new streaming buffer with validated config
func NewBuffer(config Config, w io.Writer) (*Buffer, error) {
if err := config.Validate(); err != nil {
return nil, fmt.Errorf("invalid config: %w", err)
}
if w == nil {
return nil, errors.New("writer cannot be nil")
}
return &Buffer{
config: config,
writer: w,
accum: make([]byte, 0, 4096), // 4 KiB standard initial buffer
inFence: false,
offset: -1,
}, nil
}
// calculateFenceState determines if we're inside a code fence for given data
// When resetFirst is true, it clears state before calculation (used after splits)
func (b *Buffer) calculateFenceState(data []byte, resetFirst bool) {
if resetFirst {
b.inFence = false
b.postFence = false
}
lines := bytes.Split(data, []byte("\n"))
position := 0
for _, line := range lines {
trimmed := bytes.TrimSpace(line)
// Check for code fence markers (``` or ~~~)
// Accept ``` or ~~~ optionally followed by language specifier
if (len(trimmed) >= 3 && bytes.HasPrefix(trimmed, []byte("```"))) ||
(len(trimmed) >= 3 && bytes.HasPrefix(trimmed, []byte("~~~"))) {
wasInFence := b.inFence
// Toggle fence state for standard triple-backtick fences
b.inFence = !b.inFence
// Track when we exit a fence
if wasInFence && !b.inFence {
b.postFence = true
b.lastFence = position
} else {
b.postFence = false
}
}
position += len(line) + 1 // +1 for newline
}
}
// detectFrontmatter inspects the accumulator for YAML frontmatter and records where it ends.
func (b *Buffer) detectFrontmatter() {
if b.offset != -1 {
return
}
data := b.accum
if len(data) == 0 {
return
}
// First byte must be '-'
if data[0] != '-' {
b.offset = 0
return
}
if len(data) < 3 {
// Need at least "---"
return
}
if !bytes.HasPrefix(data, []byte("---")) {
b.offset = 0
return
}
if len(data) < 4 {
return
}
newlineIdx := 3
if data[3] == '\r' {
if len(data) < 5 || data[4] != '\n' {
b.offset = 0
return
}
newlineIdx = 4
} else if data[3] != '\n' {
b.offset = 0
return
}
searchStart := newlineIdx + 1
if idx := bytes.Index(data[searchStart:], []byte("\n---\n")); idx >= 0 {
b.offset = searchStart + idx + 5
return
}
if idx := bytes.Index(data[searchStart:], []byte("\n---")); idx >= 0 && searchStart+idx+4 == len(data) {
b.offset = len(data)
return
}
if int64(len(data)) >= Windowed {
b.offset = 0
}
}
// processContentSlice updates fence state and line length tracking for new content.
func (b *Buffer) processContentSlice(content []byte, currentLineLen *int) {
if len(content) == 0 {
return
}
b.calculateFenceState(content, false)
if idx := bytes.LastIndexByte(content, '\n'); idx >= 0 {
*currentLineLen = len(content) - idx - 1
} else {
*currentLineLen += len(content)
}
if *currentLineLen > Windowed {
b.accum = append(b.accum, '\n')
*currentLineLen = 0
}
}
// handleData ingests new data into the accumulator, resolving frontmatter when possible.
func (b *Buffer) handleData(data []byte, currentLineLen *int) {
if len(data) == 0 {
return
}
// Buffer overflow protection
if int64(len(b.accum))+int64(len(data)) > Windowed {
// Force immediate flush to prevent buffer overflow
if len(b.accum) > 0 {
b.flushToSafeBoundary(true)
}
// If still too large after flush, truncate input
if int64(len(b.accum))+int64(len(data)) > Windowed {
maxAppend := Windowed - int64(len(b.accum))
if maxAppend > 0 {
data = data[:maxAppend]
} else {
return // Skip this data entirely
}
}
}
originalLen := len(b.accum)
b.accum = append(b.accum, data...)
if b.written == 0 && b.offset == -1 {
b.detectFrontmatter()
if b.offset > 0 && b.offset <= len(b.accum) {
consumed := b.offset
b.accum = b.accum[consumed:]
if len(b.accum) > 0 {
b.calculateFenceState(b.accum, true)
}
originalLen -= consumed
if originalLen < 0 {
originalLen = 0
}
b.offset = 0
} else if b.offset == 0 {
// No frontmatter detected
} else {
// Still buffering potential frontmatter
return
}
}
if len(b.accum) == 0 {
return
}
if originalLen > len(b.accum) {
originalLen = 0
}
content := b.accum[originalLen:]
b.processContentSlice(content, currentLineLen)
}
// Process reads markdown from r and streams rendered output via Output callback
func (b *Buffer) Process(ctx context.Context, r io.Reader) error {
if ctx.Err() != nil {
// Distinguish between different cancellation types:
// - context.Canceled: clean user cancellation → return nil
// - context.DeadlineExceeded: timeout → return error for proper exit code
if ctx.Err() == context.DeadlineExceeded {
return ctx.Err() // Timeout needs error for exit code
}
if ctx.Err() == context.Canceled {
return nil // Clean cancellation
}
return ctx.Err()
}
// Create cancellable context for this process
processCtx, processCancel := context.WithCancel(ctx)
defer processCancel() // Ensure goroutine cleanup
// Channel for read results
type readResult struct {
data []byte
err error
}
// Buffer the channel to avoid blocking on send
readChan := make(chan readResult, 1)
currentLineLen := 0
// Start goroutine to read from input
go func() {
defer close(readChan)
// Simple read buffer
buf := make([]byte, 4096) // 4 KiB standard buffer
for {
// Check for context cancellation before reading
select {
case <-processCtx.Done():
return
default:
}
n, err := r.Read(buf)
if n > 0 {
// Send copy of data to avoid race conditions
data := make([]byte, n)
copy(data, buf[:n])
// Try to send with context check
select {
case <-processCtx.Done():
return
case readChan <- readResult{data: data, err: err}:
}
}
// If we have an error (including EOF), send it
// Note: n > 0 with err != nil is valid (e.g., last read before EOF)
if err != nil {
// If we already sent data with the error above, don't send again
if n == 0 {
select {
case <-processCtx.Done():
return
case readChan <- readResult{err: err}:
}
}
return // Exit loop on any error including EOF
}
}
}()
// Main processing loop
for {
select {
case <-processCtx.Done():
// Context cancelled - flush accumulated content
// Tests expect partial output to be rendered on interrupt
// Use flush(false) to avoid glamour's EOF double-newline
if len(b.accum) > 0 {
if err := b.flush(); err != nil {
return err
}
}
if err := b.flushPendingPartial(); err != nil {
return err
}
return nil
case result, ok := <-readChan:
if !ok {
// Channel closed unexpectedly
if len(b.accum) > 0 {
return b.flush()
}
return nil
}
// Process the data
if len(result.data) > 0 {
b.handleData(result.data, &currentLineLen)
if b.written == 0 && b.offset == -1 {
continue
}
bufferSize := int64(len(b.accum))
// Force flush if we exceed hard limit
// This is a safety net to prevent OOM from unbounded input
// In normal operation, shouldFlush() and flushToSafeBoundary() handle this
if bufferSize > Windowed {
b.flushToSafeBoundary(true)
} else if b.shouldFlush() {
// Force flush if we've accumulated way too much
force := b.config.Window > Buffered && bufferSize > b.config.Window*2
if err := b.flushToSafeBoundary(force); err != nil {
return fmt.Errorf("opportunistic flush failed: %w", err)
}
}
// Hard limit enforcement - split pathological input at safe boundaries
if len(b.accum) >= Windowed {
// Split at safe chunk boundary to avoid glamour limits
if err := b.splitAndFlush(); err != nil {
return fmt.Errorf("pathological split failed: %w", err)
}
}
}
// Handle errors
if result.err == io.EOF {
if b.offset == -1 && b.written == 0 && len(b.accum) > 0 {
b.offset = 0
b.processContentSlice(b.accum, &currentLineLen)
}
// Final flush of all remaining content
if len(b.accum) > 0 {
return b.flushFinal()
}
// Flush any pending glamour blanks if we ended exactly on a boundary
if err := b.flushPendingFinal(); err != nil {
return err
}
// Empty input special case - glamour outputs \n\n for empty input
if b.written == 0 {
// Never had any output, render empty string to match glamour
return b.renderFinal([]byte{})
}
return nil
}
if result.err != nil {
return result.err
}
}
}
}
// shouldFlush determines if buffer should be flushed based on flow config
func (b *Buffer) shouldFlush() bool {
if b.offset == -1 && b.written == 0 {
return false
}
if b.config.Window == Buffered {
// No flow: never flush until EOF
return false
}
if b.config.Window < Buffered {
// Aggressive flow: flush at ANY newline
return bytes.Contains(b.accum, []byte("\n"))
}
// Bounded flow: flush when exceeding size threshold
return int64(len(b.accum)) >= b.config.Window
}
// findSafeBoundary finds the last safe split point
func (b *Buffer) findSafeBoundary() int {
// Never split inside code fence to preserve syntax highlighting integrity
if b.inFence {
return -1
}
// CRITICAL: If we just closed a fence, include the next line for context
if b.postFence && b.lastFence >= 0 && b.lastFence < len(b.accum) {
// Find next newline after fence to keep context together
if idx := bytes.IndexByte(b.accum[b.lastFence:], '\n'); idx >= 0 {
boundary := b.lastFence + idx + 1
// Only use this boundary if it's reasonable
if boundary < len(b.accum) {
b.postFence = false // Clear flag after use
return boundary
}
}
}
// Priority 1: Complete code blocks
if idx := b.findCodeBlockBoundary(b.accum); idx > 0 {
return idx
}
// Priority 2: Double newlines (paragraph boundaries)
if idx := bytes.LastIndex(b.accum, []byte("\n\n")); idx >= 0 {
return idx + 2
}
// Priority 3: Any newline as last resort
if idx := bytes.LastIndex(b.accum, []byte("\n")); idx >= 0 {
return idx + 1
}
// No boundary found
return -1
}
// findCodeBlockBoundary finds the end of a complete code block if present
func (b *Buffer) findCodeBlockBoundary(data []byte) int {
// Simple open/close tracking (markdown doesn't support nested fences)
inFence := false
lastCompleteEnd := -1
idx := 0
for idx < len(data) {
// Check if we're at start of line
if idx == 0 || data[idx-1] == '\n' {
// Check for fence marker (``` or ~~~)
if idx+2 < len(data) && ((data[idx] == '`' && data[idx+1] == '`' && data[idx+2] == '`') ||
(data[idx] == '~' && data[idx+1] == '~' && data[idx+2] == '~')) {
// Toggle fence state
inFence = !inFence
// Skip to end of line
for idx < len(data) && data[idx] != '\n' {
idx++
}
if idx < len(data) {
idx++ // Include newline
}
// If we just closed a fence, mark this position
if !inFence {
lastCompleteEnd = idx
}
continue
}
}
idx++
}
return lastCompleteEnd
}
// flushToSafeBoundary flushes content up to the last safe boundary
func (b *Buffer) flushToSafeBoundary(force bool) error {
if len(b.accum) == 0 {
return nil
}
// For aggressive flow (flow=-1), flush at ANY newline boundary
if b.config.Window < Buffered {
// Find the last newline in the buffer
if idx := bytes.LastIndex(b.accum, []byte("\n")); idx >= 0 {
boundary := idx + 1
toRender := b.accum[:boundary]
b.accum = b.accum[boundary:]
// Only mark as split if there's remaining content
if len(b.accum) > 0 {
b.hadSplits = true
}
err := b.render(toRender)
if err != nil {
return fmt.Errorf("render failed after %d bytes: %w", b.written, err)
}
return nil
}
// No newline found, check for pathological input
if force || len(b.accum) >= Windowed {
// Split at safe boundary to avoid glamour hanging
return b.splitAndFlush()
}
return nil
}
boundary := b.findSafeBoundary()
if boundary <= 0 {
if force {
// Forced flush: render everything
// CRITICAL: Never signal EOF for intermediate chunks!
return b.flush()
}
// If we've accumulated significantly over threshold, flush at line boundary
// Use a simple 2x threshold - if we're over 2x the window size, force a flush
if b.config.Window > Buffered && int64(len(b.accum)) > b.config.Window*2 {
// Find a line boundary as fallback
if idx := bytes.LastIndex(b.accum, []byte("\n")); idx > 0 {
boundary = idx + 1
} else {
// No line boundary either, split at safe chunk size
return b.splitAndFlush()
}
} else {
// No safe boundary: keep accumulating
return nil
}
}
// Split at boundary
toRender := b.accum[:boundary]
b.accum = b.accum[boundary:]
// Only mark as split if there's remaining content (not just flushing everything)
if len(b.accum) > 0 {
b.hadSplits = true
}
// Recalculate fence state for remaining content after split
// CRITICAL: Don't reset fence state - preserve it across boundaries
b.calculateFenceState(b.accum, false)
// Debug output
// Reset cache since buffer changed
// This is an intermediate flush at a boundary, never final
// Only explicit flush(true) calls should be final
err := b.render(toRender)
if err != nil {
return fmt.Errorf("split flush failed: %w", err)
}
return nil
}
// flush renders all accumulated content
func (b *Buffer) flush() error {
if len(b.accum) == 0 {
return nil
}
toRender := b.accum
b.accum = b.accum[:0] // Clear but keep capacity
// Reset cache since buffer is cleared
// Reset fence state after flush (accumulator is now empty)
b.inFence = false
return b.render(toRender)
}
// flushFinal renders all accumulated content as final output
func (b *Buffer) flushFinal() error {
if len(b.accum) == 0 {
return nil
}
toRender := b.accum
b.accum = b.accum[:0] // Clear but keep capacity
// Reset fence state after flush (accumulator is now empty)
b.inFence = false
return b.renderFinal(toRender)
}
// splitAndFlush splits pathological input at safe boundaries to avoid glamour hanging
func (b *Buffer) splitAndFlush() error {
if len(b.accum) == 0 {
return nil
}
// Split at SafeChunkSize boundary to stay under glamour's limit
for len(b.accum) >= Windowed {
// Extract a safe chunk
chunk := b.accum[:Windowed]
b.accum = b.accum[Windowed:]
b.hadSplits = true // Mark that we've split the input
// Recalculate fence state for remaining content after split
// CRITICAL: Don't reset fence state - preserve it across boundaries
b.calculateFenceState(b.accum, false)
// Render the chunk
if err := b.render(chunk); err != nil {
return fmt.Errorf("flush render failed: %w", err)
}
}
// If there's remaining data less than chunk size, flush it
if len(b.accum) > 0 {
return b.flush()
}
return nil
}
const glamourBlankSequence = "\n \n"
const glamourBlankLen = len(glamourBlankSequence)
var glamourBlankBytes = []byte(glamourBlankSequence)
const pendingBlankSequence = " \n"
const pendingBlankLen = len(pendingBlankSequence)
var pendingBlankBytes = []byte(pendingBlankSequence)
func trailingGlamourBlanks(data []byte) int {
suffix := 0
for len(data) >= glamourBlankLen && bytes.HasSuffix(data, glamourBlankBytes) {
suffix += glamourBlankLen
data = data[:len(data)-glamourBlankLen]
}
return suffix
}
func convertGlamourBlankToFinal(data []byte) []byte {
if len(data) == 0 {
return nil
}
count := len(data) / pendingBlankLen
if count <= 0 {
return nil
}
result := make([]byte, count)
for i := range result {
result[i] = '\n'
}
return result
}
func (b *Buffer) writeChunk(output []byte, final bool) error {
// Enforce total output size limit
if b.written+int64(len(output)) > Windowed {
// Calculate how much we can still write
remaining := Windowed - b.written
if remaining <= 0 {
// Already at limit, don't write anything more
return nil
}
// Truncate output to fit within limit
output = output[:remaining]
}
swallowEPIPE := !final
if final {
suffixLen := trailingGlamourBlanks(output)
count := suffixLen / glamourBlankLen
var suffixConverted []byte
if count > 0 {
suffixStart := len(output) - suffixLen
output = append(output[:suffixStart], bytes.Repeat([]byte("\n"), count)...)
suffixConverted = make([]byte, count)
for i := range suffixConverted {
suffixConverted[i] = '\n'
}
}
convertPending := len(output) == 0 && len(suffixConverted) == 0
if len(b.pendingBlank) > 0 {
pending := b.pendingBlank
if convertPending {
pending = convertGlamourBlankToFinal(pending)
}
if err := b.writeToWriter(pending, swallowEPIPE); err != nil {
return err
}
b.pendingBlank = b.pendingBlank[:0]
}
if len(output) > 0 {
if err := b.writeToWriter(output, swallowEPIPE); err != nil {
return err
}
}
if len(suffixConverted) > 0 {
return b.writeToWriter(suffixConverted, swallowEPIPE)
}
return nil
}
// Non-final chunk
if len(b.pendingBlank) > 0 {
if err := b.writeToWriter(b.pendingBlank, swallowEPIPE); err != nil {
return err
}
b.pendingBlank = b.pendingBlank[:0]
}
suffixLen := trailingGlamourBlanks(output)
if suffixLen > 0 {
count := suffixLen / glamourBlankLen
suffixStart := len(output) - suffixLen
output = append(output[:suffixStart], bytes.Repeat([]byte("\n"), count)...)
b.pendingBlank = b.pendingBlank[:0]
for i := 0; i < count; i++ {
b.pendingBlank = append(b.pendingBlank, pendingBlankBytes...)
}
}
if len(output) > 0 {
return b.writeToWriter(output, swallowEPIPE)
}
return nil
}
func (b *Buffer) writeToWriter(data []byte, swallowEPIPE bool) error {
if len(data) == 0 {
return nil
}
n, err := b.writer.Write(data)
if err != nil {
if swallowEPIPE && errors.Is(err, syscall.EPIPE) {
b.written += int64(n)
return nil
}
return fmt.Errorf("write failed after %d bytes: %w", b.written, err)
}
b.written += int64(n)
return nil
}
func (b *Buffer) flushPendingFinal() error {
if len(b.pendingBlank) == 0 {
return nil
}
converted := convertGlamourBlankToFinal(b.pendingBlank)
b.pendingBlank = b.pendingBlank[:0]
return b.writeToWriter(converted, false)
}
func (b *Buffer) flushPendingPartial() error {
if len(b.pendingBlank) == 0 {
return nil
}
count := len(b.pendingBlank) / pendingBlankLen
newline := make([]byte, count)
for i := range newline {
newline[i] = '\n'
}
b.pendingBlank = b.pendingBlank[:0]
return b.writeToWriter(newline, true)
}
// renderFinal is like render but for the final output
func (b *Buffer) renderFinal(data []byte) error {
output, err := b.config.Render(data)
if err != nil {
return err
}
// Check if this is identity renderer (returns exactly what was given)
isIdentity := bytes.Equal(output, data)
// Conservative glamour detection - only apply spacing fixes to actual glamour output
// Glamour output has consistent " " prefix on content lines
isLikelyGlamour := !isIdentity &&
len(output) > 3 &&
bytes.HasPrefix(output, []byte("\n ")) &&
bytes.Contains(output, []byte("\n ")) // Has multiple glamour-formatted lines
// Simple blank line handling for glamour output only
// For the final render, we need special handling of the trailing newlines
if isLikelyGlamour {
// Special handling: glamour ends with \n\n, we should preserve that
endsWithDoubleNewline := bytes.HasSuffix(output, []byte("\n\n"))
// Do the normal blank line replacement
output = bytes.ReplaceAll(output, []byte("\n\n"), []byte("\n \n"))
// If it originally ended with \n\n, restore that at the end
if endsWithDoubleNewline && bytes.HasSuffix(output, []byte("\n \n")) {
// Remove the spaces from the final empty line
output = output[:len(output)-4] // Remove " \n"
output = append(output, []byte("\n\n")...) // Add back "\n\n"
}
}
// Strip leading newline from non-first chunks (glamour only)
if isLikelyGlamour && b.written > 0 && len(output) > 0 && output[0] == '\n' {
output = output[1:]
}
// Skip completely empty output
if len(output) == 0 {
return b.writeChunk(output, true)
}
return b.writeChunk(output, true)
}
// render processes markdown through glamour and outputs directly to writer
func (b *Buffer) render(data []byte) error {
output, err := b.config.Render(data)
if err != nil {
return err
}
// Check if this is identity renderer (returns exactly what was given)
isIdentity := bytes.Equal(output, data)
// Conservative glamour detection - only apply spacing fixes to actual glamour output
// Glamour output has consistent " " prefix on content lines
isLikelyGlamour := !isIdentity &&
len(output) > 3 &&
bytes.HasPrefix(output, []byte("\n ")) &&
bytes.Contains(output, []byte("\n ")) // Has multiple glamour-formatted lines
// Simple blank line handling for glamour output only
// Replace \n\n with \n \n to maintain spacing
if isLikelyGlamour {
output = bytes.ReplaceAll(output, []byte("\n\n"), []byte("\n \n"))
// Ensure trailing empty lines have proper spacing
// This ensures consistency with glow.orig output
if bytes.HasSuffix(output, []byte("\n \n")) {
// Already has proper spacing
} else if bytes.HasSuffix(output, []byte("\n\n")) {
// Replace final empty line with properly spaced line
output = bytes.TrimSuffix(output, []byte("\n"))
output = append(output, []byte(" \n")...)
}
}
// Strip leading newline from non-first chunks (glamour only)
if isLikelyGlamour && b.written > 0 && len(output) > 0 && output[0] == '\n' {
output = output[1:]
}
return b.writeChunk(output, false)
}
// Flow is the main entry point for streaming markdown rendering with configurable flow control.
//
// Flow processes markdown content from reader r and writes rendered output to writer w using the provided render function.
// The window parameter controls buffering behavior:
//
// - window == -1 (Unbuffered): Aggressive flushing at newline boundaries for minimal latency
// - window == 0 (Buffered): Maximum buffering until EOF for optimal throughput
// - window > 0 (Windowed): Bounded buffering with specified byte limit for balanced performance
//
// The render function should accept markdown bytes and return rendered output (e.g., HTML or styled text).
// Flow handles smart boundary detection to avoid breaking markdown structures like code fences.
//
// Returns an error if validation fails or processing encounters unrecoverable issues.
func Flow(ctx context.Context, r io.Reader, w io.Writer, window int64, render RenderFunc) error {
// Validate inputs
switch {
case ctx == nil:
return fmt.Errorf("context cannot be nil")
case r == nil:
return fmt.Errorf("input reader cannot be nil")
case w == nil:
return fmt.Errorf("output writer cannot be nil")
case render == nil:
return fmt.Errorf("render function cannot be nil")
case window < Unbuffered || (window > 0 && window > Windowed):
return fmt.Errorf("window must be -1 (block), 0 (stream), or 1-%d (bytes)", Windowed)
}
// Create buffer configuration
config := Config{
Window: window,
Render: render,
}
// Create buffer with validation
buffer, err := NewBuffer(config, w)
if err != nil {
return fmt.Errorf("failed to create buffer: %w", err)
}
// Process the stream
return buffer.Process(ctx, r)
}
+193
View File
@@ -0,0 +1,193 @@
package markdown
import (
"bytes"
"context"
"fmt"
"os"
"strings"
"sync"
"github.com/charmbracelet/glamour"
)
// StreamRenderer handles ChatGPT-style hybrid streaming:
// 1. Shows plain text immediately (character-by-character)
// 2. Freezes completed blocks with styling (no jumping)
// 3. Only streams new plain text for incomplete blocks
type StreamRenderer struct {
mu sync.Mutex
glamour *glamour.TermRenderer
ctx context.Context
cancel context.CancelFunc
frozenOutput string // Completed, styled blocks (immutable)
currentBlock bytes.Buffer // Current block being streamed (plain text)
allText bytes.Buffer // Full accumulated text for context
plainCharsWritten int // Number of plain chars displayed in current block
inputDebug *os.File
}
// NewStreamRenderer creates a new streaming markdown renderer
func NewStreamRenderer() (*StreamRenderer, error) {
// Create Glamour renderer
renderer, err := glamour.NewTermRenderer(
glamour.WithStylePath("dark"),
glamour.WithWordWrap(80),
)
if err != nil {
return nil, fmt.Errorf("failed to create glamour renderer: %w", err)
}
ctx, cancel := context.WithCancel(context.Background())
sr := &StreamRenderer{
glamour: renderer,
ctx: ctx,
cancel: cancel,
}
// Enable debug logging if CLINE_MARKDOWN_DEBUG env var is set
if os.Getenv("CLINE_MARKDOWN_DEBUG") != "" {
debugPath := os.Getenv("CLINE_MARKDOWN_DEBUG")
if inputFile, err := os.Create(debugPath + ".input.md"); err == nil {
sr.inputDebug = inputFile
}
}
return sr, nil
}
// WriteIncremental implements ChatGPT-style block-based streaming:
// 1. Show plain text immediately
// 2. Freeze completed blocks (no re-rendering)
// 3. Only re-render current incomplete block
func (sr *StreamRenderer) WriteIncremental(text string) error {
sr.mu.Lock()
defer sr.mu.Unlock()
// Log input
if sr.inputDebug != nil {
sr.inputDebug.WriteString(text)
}
// Append to all text for full context
sr.allText.WriteString(text)
// Append to current block
sr.currentBlock.WriteString(text)
// Show plain text immediately (ChatGPT-style)
fmt.Print(text)
sr.plainCharsWritten += len(text)
// Check if we should freeze this block and re-render
if sr.shouldFreezeBlock(text) {
sr.freezeCurrentBlock()
}
return nil
}
// shouldFreezeBlock determines if we should freeze the current block
// This happens at logical boundaries like paragraph breaks
func (sr *StreamRenderer) shouldFreezeBlock(newText string) bool {
currentBlockText := sr.currentBlock.String()
// Freeze on double newline (paragraph break)
if strings.HasSuffix(currentBlockText, "\n\n") {
return true
}
// Freeze on closing code block
if strings.HasSuffix(currentBlockText, "```\n") {
// Count code fences in current block
fenceCount := strings.Count(currentBlockText, "```")
if fenceCount >= 2 && fenceCount%2 == 0 {
return true
}
}
// Freeze every ~5 lines to keep blocks manageable
lineCount := strings.Count(currentBlockText, "\n")
if lineCount >= 5 {
return true
}
return false
}
// freezeCurrentBlock renders the current block and adds it to frozen output
func (sr *StreamRenderer) freezeCurrentBlock() {
if sr.currentBlock.Len() == 0 {
return
}
// Clear the plain text we just printed
plainText := sr.currentBlock.String()
lineCount := strings.Count(plainText, "\n")
if lineCount > 0 {
// Move cursor up to start of current block
fmt.Printf("\033[%dA", lineCount)
// Clear from cursor down
fmt.Print("\033[J")
} else {
// Clear current line
fmt.Print("\r\033[K")
}
// Render the full context (frozen + current) to get correct styling
fullText := sr.frozenOutput + sr.currentBlock.String()
rendered, err := sr.glamour.Render(fullText)
if err != nil {
// Fallback: keep plain text
fmt.Print(sr.frozenOutput)
fmt.Print(plainText)
sr.currentBlock.Reset()
sr.plainCharsWritten = 0
return
}
// Print the fully rendered output
fmt.Print(rendered)
// Update frozen output to include this block
sr.frozenOutput = fullText
// Reset current block
sr.currentBlock.Reset()
sr.plainCharsWritten = 0
}
// FlushMessage completes the current message with final markdown render
func (sr *StreamRenderer) FlushMessage() error {
sr.mu.Lock()
defer sr.mu.Unlock()
// Freeze any remaining content in current block
if sr.currentBlock.Len() > 0 {
sr.freezeCurrentBlock()
}
// Add separator
fmt.Println()
// Reset everything
sr.frozenOutput = ""
sr.currentBlock.Reset()
sr.allText.Reset()
sr.plainCharsWritten = 0
return nil
}
// Close shuts down the streaming renderer gracefully
func (sr *StreamRenderer) Close() error {
sr.cancel()
if sr.inputDebug != nil {
sr.inputDebug.Close()
}
return nil
}
+1 -1
View File
@@ -9,6 +9,7 @@ github.com/envoyproxy/go-control-plane/ratelimit v0.1.0/go.mod h1:Wk+tMFAFbCXaJP
github.com/envoyproxy/protoc-gen-validate v1.2.1/go.mod h1:d/C80l/jxXLdfEIhX1W2TmLfsJ31lvEjwamM4DxlWXU=
github.com/go-jose/go-jose/v4 v4.1.1/go.mod h1:BdsZGqgdO3b6tTc6LSE56wcDbMMLuPsw5d4ZD5f94kA=
github.com/golang/glog v1.2.5/go.mod h1:6AhwSGph0fcJtXVM/PEHPqZlFeoLxhs7/t5UDAwmO+w=
github.com/hexops/gotextdiff v1.0.3/go.mod h1:pSWU5MAI3yDq+fZBTazCSJysOMbxWL1BSow5/V2vxeg=
github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10/go.mod h1:t/avpk3KcrXxUnYOhZhMXJlSEyie6gQbtLq5NM3loB8=
github.com/spiffe/go-spiffe/v2 v2.5.0/go.mod h1:P+NxobPc6wXhVtINNtFjNWGBTreew1GBUCwT2wPmb7g=
github.com/zeebo/errs v1.4.0/go.mod h1:sgbWHsvVuTPHcqJJGQ1WhI5KbWlHYz+2+2C/LSEtCw4=
@@ -17,7 +18,6 @@ golang.org/x/crypto v0.39.0/go.mod h1:L+Xg3Wf6HoL4Bn4238Z6ft6KfEpN0tJGo53AAPC632
golang.org/x/mod v0.25.0/go.mod h1:IXM97Txy2VM4PJ3gI61r1YEk/gAj6zAHN3AdZt6S9Ww=
golang.org/x/oauth2 v0.30.0/go.mod h1:B++QgG3ZKulg6sRPGD/mqlHQs5rB3Ml9erfeDY7xKlU=
golang.org/x/sync v0.15.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/term v0.32.0/go.mod h1:uZG1FhGx848Sqfsq4/DlJr3xGGsYMu/L5GW4abiaEPQ=
golang.org/x/tools v0.33.0/go.mod h1:CIJMaWEY88juyUfo7UbgPqbC8rU2OqfAV1h2Qp0oMYI=
golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
google.golang.org/genproto/googleapis/api v0.0.0-20250707201910-8d1bb00bc6a7/go.mod h1:kXqgZtrWaf6qS3jZOCnCH7WYfrvFjkC51bM8fz3RsCA=