pico
created pr with
ps-190
changed status to
accepted
cmds
checkout latest patchset:
ssh pr.pico.sh print pr-104 | git am -3checkout any patchset in a patch request:
ssh pr.pico.sh print ps-X | git am -3add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 104set PR to open (enables RSS notifications):
ssh pr.pico.sh pr open 104set PR to draft (stops RSS notifications):
ssh pr.pico.sh pr draft 104
Patchset
ps-190
feat(pipe): ?prefix= query param for pub and pipe
Eric Bower
2026-01-22T23:51:10ZSemantic diff summary
0 added,
3 modified,
0 signature changed,
0 removed
across 1 analyzed file
+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
+ }
}
}
}()