pico

created pr with 103.1 on 2026-01-05T14:52:02Z · by c8ef7d19
added 103.2 on 2026-01-08T01:19:34Z · by c8ef7d19
1: 66cf9b9 = 1: 66cf9b9 feat(pipe): add pipe_monitors table
2: 83b0c70 = 2: 83b0c70 chore(pipe): add new db methods
3: 43aa4a1 = 3: 43aa4a1 chore(pipe): add db impl
4: a75186b = 4: a75186b chore(pipe): add tests for monitoring
5: 9c7fcda = 5: 9c7fcda feat(pipe): add monitor cli
6: 2b86dec = 6: 2b86dec feat(pipe): status and rss commands
7: 8df2aa2 = 7: 8df2aa2 feat(pipe): monitor pub and pipe cmd
8: 11a5bce = 8: 11a5bce chore(pipe): monitor help text
9: 670350c = 9: 670350c refactor(pipe): monitor pipes on throttle interval
10: ade2b83 = 10: ade2b83 fix: window and ping fixes
11: 68f5cd1 = 11: 68f5cd1 chore(pipe): add tests
-: ------- > 12: 74fb40f chore(pipe): add logging stmts
-: ------- > 13: c422818 refactor(pipe): create pipe monitors history table
-: ------- > 14: e0c3a88 fix(pipe): only update monitor if within window
-: ------- > 15: 9a8fbde chore(pipe): implement monitor history db interface
-: ------- > 16: f168507 feat(pipe): record historical monitors
-: ------- > 17: 16df77f feat(pipe): monitor calculate uptime
-: ------- > 18: 10ffdf1 feat(pipe): uptime cli command
-: ------- > 19: 8231ac5 chore(pipe): monitor history stubs
-: ------- > 20: 65a63a9 fix(pipe): uptime cmd
added 103.3 on 2026-02-01T17:16:24Z · by c8ef7d19
3: 43aa4a1 ! 1: 603ca6e chore(pubsub): add more tests
1: 66cf9b9 < -: ------- feat(pipe): add pipe_monitors table
2: 83b0c70 < -: ------- chore(pipe): add new db methods
4: a75186b ! 2: 17e00b2 feat(pubsub): round robin
13: c422818 ! 3: e3136bd fix(pubsub): check for eof before processing and skip empty byte reads
-: ------- > 4: 9a6d19e fix: rr
5: 9c7fcda < -: ------- feat(pipe): add monitor cli
16: f168507 ! 5: 5b3f3a1 fix: sending 0 byte read
6: 2b86dec < -: ------- feat(pipe): status and rss commands
11: 68f5cd1 ! 6: 4fff471 refactor: fixes
7: 8df2aa2 < -: ------- feat(pipe): monitor pub and pipe cmd
14: e0c3a88 ! 7: d0dfc85 chore: SetDispatch on Broker
8: 11a5bce < -: ------- chore(pipe): monitor help text
9: 670350c < -: ------- refactor(pipe): monitor pipes on throttle interval
10: ade2b83 < -: ------- fix: window and ping fixes
12: 74fb40f < -: ------- chore(pipe): add logging stmts
15: 9a8fbde < -: ------- chore(pipe): implement monitor history db interface
17: 16df77f < -: ------- feat(pipe): monitor calculate uptime
18: 10ffdf1 < -: ------- feat(pipe): uptime cli command
19: 8231ac5 < -: ------- chore(pipe): monitor history stubs
20: 65a63a9 < -: ------- fix(pipe): uptime cmd
cmds
checkout latest patchset:
ssh pr.pico.sh print 103 | git am -3
checkout any patchset in a patch request:
ssh pr.pico.sh print 103.[rev] | git am -3
add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 103

Patchset 103.2 on 2026-01-08T01:19:34Z · commit 65a63a9

+27 -21 pkg/apps/pipe/cli.go #
......@@ -3,6 +3,8 @@ package pipe
33 import (
44 "bytes"
55 "context"
6+ "database/sql"
7+ "errors"
68 "flag"
79 "fmt"
810 "io"
......@@ -90,13 +92,6 @@ func Middleware(handler *CliHandler) pssh.SSHServerMiddleware {
9092 sesh.Fatal(err)
9193 }
9294 return next(sesh)
93- case "uptime":
94- err := handler.uptime(cliCmd, user)
95- if err != nil {
96- logger.Error("uptime cmd", "err", err)
97- sesh.Fatal(err)
98- }
99- return next(sesh)
10095 case "rss":
10196 err := handler.rss(cliCmd, user)
10297 if err != nil {
......@@ -170,6 +165,12 @@ func Middleware(handler *CliHandler) pssh.SSHServerMiddleware {
170165 logger.Error("pipe cmd", "err", err)
171166 sesh.Fatal(err)
172167 }
168+ case "uptime":
169+ err := handler.uptime(cliCmd, topic, user)
170+ if err != nil {
171+ logger.Error("uptime cmd", "err", err)
172+ sesh.Fatal(err)
173+ }
173174 }
174175
175176 return next(sesh)
......@@ -449,11 +450,18 @@ func (handler *CliHandler) status(cmd *CliCmd, user *db.User) error {
449450 return nil
450451 }
451452
452-func (handler *CliHandler) uptime(cmd *CliCmd, user *db.User) error {
453+func (handler *CliHandler) uptime(cmd *CliCmd, topic string, user *db.User) error {
453454 if user == nil {
454455 return fmt.Errorf("access denied")
455456 }
456457
458+ if topic == "" {
459+ _, _ = fmt.Fprintln(cmd.sesh, "usage: uptime <topic> [--from <time>] [--to <time>]")
460+ _, _ = fmt.Fprintln(cmd.sesh, " --from: start time (RFC3339 or duration like '24h', '7d', default: 24h)")
461+ _, _ = fmt.Fprintln(cmd.sesh, " --to: end time (RFC3339, default: now)")
462+ return nil
463+ }
464+
457465 fs := flag.NewFlagSet("uptime", flag.ContinueOnError)
458466 fs.SetOutput(cmd.sesh)
459467 fromStr := fs.String("from", "", "start time (RFC3339 or duration like '24h', '7d')")
......@@ -463,23 +471,21 @@ func (handler *CliHandler) uptime(cmd *CliCmd, user *db.User) error {
463471 return nil
464472 }
465473
466- args := fs.Args()
467- if len(args) == 0 {
468- _, _ = fmt.Fprintln(cmd.sesh, "usage: uptime <topic> [--from <time>] [--to <time>]")
469- _, _ = fmt.Fprintln(cmd.sesh, " --from: start time (RFC3339 or duration like '24h', '7d', default: 24h)")
470- _, _ = fmt.Fprintln(cmd.sesh, " --to: end time (RFC3339, default: now)")
471- return nil
472- }
473-
474- topic := args[0]
474+ topicResult := resolveTopic(TopicResolveInput{
475+ UserName: cmd.userName,
476+ Topic: topic,
477+ IsAdmin: cmd.isAdmin,
478+ IsPublic: false,
479+ })
480+ resolvedTopic := topicResult.Name
475481
476- monitor, err := handler.DBPool.FindPipeMonitorByTopic(user.ID, topic)
482+ monitor, err := handler.DBPool.FindPipeMonitorByTopic(user.ID, resolvedTopic)
477483 if err != nil {
484+ if errors.Is(err, sql.ErrNoRows) {
485+ return fmt.Errorf("monitor not found: %s", topic)
486+ }
478487 return fmt.Errorf("failed to find monitor: %w", err)
479488 }
480- if monitor == nil {
481- return fmt.Errorf("monitor not found: %s", topic)
482- }
483489
484490 now := time.Now().UTC()
485491 to := now
+3 -2 pkg/db/postgres/storage.go #
......@@ -1567,9 +1567,10 @@ func (me *PsqlDB) FindPipeMonitorsByUser(userID string) ([]*db.PipeMonitor, erro
15671567 }
15681568
15691569 func (me *PsqlDB) InsertPipeMonitorHistory(monitorID string, windowDur time.Duration, windowEnd, lastPing *time.Time) error {
1570+ durStr := fmt.Sprintf("%d seconds", int64(windowDur.Seconds()))
15701571 _, err := me.Db.Exec(
1571- `INSERT INTO pipe_monitors_history (monitor_id, window_dur, window_end, last_ping) VALUES ($1, $2, $3, $4)`,
1572- monitorID, windowDur, windowEnd, lastPing,
1572+ `INSERT INTO pipe_monitors_history (monitor_id, window_dur, window_end, last_ping) VALUES ($1, $2::interval, $3, $4)`,
1573+ monitorID, durStr, windowEnd, lastPing,
15731574 )
15741575 return err
15751576 }
Back to top