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 83b0c70

+36 -0 pkg/db/db.go #
......@@ -5,6 +5,7 @@ import (
55 "database/sql/driver"
66 "encoding/json"
77 "errors"
8+ "fmt"
89 "regexp"
910 "time"
1011 )
......@@ -378,6 +379,36 @@ type TunsEventLog struct {
378379 CreatedAt *time.Time `json:"created_at" db:"created_at"`
379380 }
380381
382+type PipeMonitor struct {
383+ ID string `json:"id" db:"id"`
384+ UserId string `json:"user_id" db:"user_id"`
385+ Topic string `json:"topic" db:"topic"`
386+ WindowDur time.Duration `json:"window_dur" db:"window_dur"`
387+ WindowEnd *time.Time `json:"window_end" db:"window_end"`
388+ LastPing *time.Time `json:"last_ping" db:"last_ping"`
389+ CreatedAt *time.Time `json:"created_at" db:"created_at"`
390+ UpdatedAt *time.Time `json:"updated_at" db:"updated_at"`
391+}
392+
393+func (m *PipeMonitor) Status() error {
394+ windowStart := m.WindowEnd.Add(-m.WindowDur)
395+ lastPingAfterStart := m.LastPing.After(windowStart)
396+ if !lastPingAfterStart {
397+ return fmt.Errorf("last ping before window start: last_ping=%s window_start=%s", m.LastPing, windowStart)
398+ }
399+ lastPingBeforeEnd := m.LastPing.Before(*m.WindowEnd)
400+ if !lastPingBeforeEnd {
401+ // should not happen but just for data validity we add it
402+ return fmt.Errorf("last ping after window end: last_ping=%s window_end=%s", m.LastPing, m.WindowEnd)
403+ }
404+ return nil
405+}
406+
407+func (m *PipeMonitor) GetNextWindow() *time.Time {
408+ win := m.WindowEnd.Add(m.WindowDur)
409+ return &win
410+}
411+
381412 var NameValidator = regexp.MustCompile("^[a-zA-Z0-9]{1,50}$")
382413 var DenyList = []string{
383414 "admin",
......@@ -468,5 +499,10 @@ type DB interface {
468499 FindPubkeysInAccessLogs(userID string) ([]string, error)
469500 FindAccessLogsByPubkey(pubkey string, fromDate *time.Time) ([]*AccessLog, error)
470501
502+ UpsertPipeMonitor(userID, topic string, dur time.Duration, winEnd *time.Time) error
503+ UpdatePipeMonitorLastPing(userID, topic string, lastPing *time.Time) error
504+ RemovePipeMonitor(userID, topic string) error
505+ FindPipeMonitorByTopic(userID, topic string) (*PipeMonitor, error)
506+
471507 Close() error
472508 }
Back to top