pico

created pr with 104.1 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 104 | git am -3
checkout any patchset in a patch request:
ssh pr.pico.sh print 104.[rev] | 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

Patchset 104.1 on 2026-01-22T23:52:32Z · commit 633c17a

+53 -7 pkg/apps/pipe/api.go #
......@@ -2,6 +2,7 @@ package pipe
22
33 import (
44 "bufio"
5+ "bytes"
56 "context"
67 "errors"
78 "fmt"
......@@ -158,6 +159,8 @@ func handlePub(pubsub bool) http.HandlerFunc {
158159 params += fmt.Sprintf(" -a=%s", cleanList)
159160 }
160161
162+ prefix := r.URL.Query().Get("prefix")
163+
161164 var wg sync.WaitGroup
162165
163166 reader := bufio.NewReaderSize(r.Body, 1)
......@@ -245,7 +248,12 @@ func handlePub(pubsub bool) http.HandlerFunc {
245248 case <-r.Context().Done():
246249 break outer
247250 default:
248- n, err := p.Write(first)
251+ messageToWrite := first
252+ if prefix != "" {
253+ messageToWrite = append([]byte(prefix), messageToWrite...)
254+ }
255+
256+ n, err := p.Write(messageToWrite)
249257 if err != nil {
250258 logger.Error("pub write error", "topic", topic, "info", clientInfo, "err", err.Error())
251259 http.Error(w, "server error", http.StatusInternalServerError)
......@@ -314,6 +322,8 @@ func handlePipe() http.HandlerFunc {
314322 params += fmt.Sprintf(" -a=%s", cleanList)
315323 }
316324
325+ prefix := r.URL.Query().Get("prefix")
326+
317327 id := uuid.NewString()
318328
319329 p, err := sshClient.AddSession(id, fmt.Sprintf("pipe %s %s", params, topic), 0, -1, -1)
......@@ -364,6 +374,8 @@ func handlePipe() http.HandlerFunc {
364374 wg.Done()
365375 }()
366376
377+ var messageBuffer []byte
378+
367379 for {
368380 buf := make([]byte, 32*1024)
369381
......@@ -373,12 +385,46 @@ func handlePipe() http.HandlerFunc {
373385 break
374386 }
375387
376- buf = buf[:n]
377-
378- err = c.WriteMessage(messageType, buf)
379- if err != nil {
380- logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
381- break
388+ messageBuffer = append(messageBuffer, buf[:n]...)
389+
390+ if prefix != "" {
391+ // Buffer and split on prefix boundaries
392+ for {
393+ firstIdx := bytes.Index(messageBuffer, []byte(prefix))
394+ if firstIdx == -1 {
395+ // No prefix found, clear buffer (shouldn't happen in normal use)
396+ messageBuffer = nil
397+ break
398+ }
399+
400+ // Look for next prefix after the first one
401+ secondIdx := bytes.Index(messageBuffer[firstIdx+len(prefix):], []byte(prefix))
402+ if secondIdx == -1 {
403+ // No complete message yet, keep buffer as is
404+ break
405+ }
406+
407+ // We have a complete message, extract and send it
408+ messageToSend := messageBuffer[firstIdx : firstIdx+len(prefix)+secondIdx]
409+ err = c.WriteMessage(messageType, messageToSend)
410+ if err != nil {
411+ logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
412+ break
413+ }
414+
415+ // Update buffer to remove sent message
416+ messageBuffer = messageBuffer[firstIdx+len(prefix)+secondIdx:]
417+ }
418+ } else {
419+ // No prefix set, send all data as-is
420+ if len(messageBuffer) > 0 {
421+ err = c.WriteMessage(messageType, messageBuffer)
422+ if err != nil {
423+ logger.Error("pipe write error", "topic", topic, "info", clientInfo, "err", err.Error())
424+ break
425+ }
426+ messageBuffer = nil
427+ }
382428 }
383429 }
384430 }()
Back to top