pico

created pr with ps-190 on 2026-01-22T23:52:32Z · by c8ef7d19
changed status to accepted on 2026-02-26T01:40:08Z · by c8ef7d19
cmds
checkout latest patchset:
ssh pr.pico.sh print pr-104 | git am -3
checkout any patchset in a patch request:
ssh pr.pico.sh print ps-X | git am -3
add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 104
set PR to open (enables RSS notifications):
ssh pr.pico.sh pr open 104
set PR to draft (stops RSS notifications):
ssh pr.pico.sh pr draft 104
+53 -7 pkg/apps/pipe/api.go #
@@ -2,6 +2,7 @@ package pipe
 
 import (
 	"bufio"
+	"bytes"
 	"context"
 	"errors"
 	"fmt"
@@ -158,6 +159,8 @@ func handlePub(pubsub bool) http.HandlerFunc {
 			params += fmt.Sprintf(" -a=%s", cleanList)
 		}
 
+		prefix := r.URL.Query().Get("prefix")
+
 		var wg sync.WaitGroup
 
 		reader := bufio.NewReaderSize(r.Body, 1)
@@ -245,7 +248,12 @@ func handlePub(pubsub bool) http.HandlerFunc {
 			case <-r.Context().Done():
 				break outer
 			default:
-				n, err := p.Write(first)
+				messageToWrite := first
+				if prefix != "" {
+					messageToWrite = append([]byte(prefix), messageToWrite...)
+				}
+
+				n, err := p.Write(messageToWrite)
 				if err != nil {
 					logger.Error("pub write error", "topic", topic, "info", clientInfo, "err", err.Error())
 					http.Error(w, "server error", http.StatusInternalServerError)
@@ -314,6 +322,8 @@ func handlePipe() http.HandlerFunc {
 			params += fmt.Sprintf(" -a=%s", cleanList)
 		}
 
+		prefix := r.URL.Query().Get("prefix")
+
 		id := uuid.NewString()
 
 		p, err := sshClient.AddSession(id, fmt.Sprintf("pipe %s %s", params, topic), 0, -1, -1)
@@ -364,6 +374,8 @@ func handlePipe() http.HandlerFunc {
 				wg.Done()
 			}()
 
+			var messageBuffer []byte
+
 			for {
 				buf := make([]byte, 32*1024)
 
@@ -373,12 +385,46 @@ func handlePipe() http.HandlerFunc {
 					break
 				}
 
-				buf = buf[:n]
-
-				err = c.WriteMessage(messageType, buf)
-				if err != nil {
-					logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
-					break
+				messageBuffer = append(messageBuffer, buf[:n]...)
+
+				if prefix != "" {
+					// Buffer and split on prefix boundaries
+					for {
+						firstIdx := bytes.Index(messageBuffer, []byte(prefix))
+						if firstIdx == -1 {
+							// No prefix found, clear buffer (shouldn't happen in normal use)
+							messageBuffer = nil
+							break
+						}
+
+						// Look for next prefix after the first one
+						secondIdx := bytes.Index(messageBuffer[firstIdx+len(prefix):], []byte(prefix))
+						if secondIdx == -1 {
+							// No complete message yet, keep buffer as is
+							break
+						}
+
+						// We have a complete message, extract and send it
+						messageToSend := messageBuffer[firstIdx : firstIdx+len(prefix)+secondIdx]
+						err = c.WriteMessage(messageType, messageToSend)
+						if err != nil {
+							logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
+							break
+						}
+
+						// Update buffer to remove sent message
+						messageBuffer = messageBuffer[firstIdx+len(prefix)+secondIdx:]
+					}
+				} else {
+					// No prefix set, send all data as-is
+					if len(messageBuffer) > 0 {
+						err = c.WriteMessage(messageType, messageBuffer)
+						if err != nil {
+							logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
+							break
+						}
+						messageBuffer = nil
+					}
 				}
 			}
 		}()
Back to top