jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360package main
import ( "bytes" "context" "encoding/json" "fmt" "net/http" "net/http/httptest" "os" "os/exec" "path/filepath" "runtime" "sort" "strconv" "strings" "sync" "sync/atomic" "syscall" "time"
js "github.com/bluesky-social/jetstream" "github.com/coder/websocket" "github.com/klauspost/compress/zstd")
// Performance is a separate lane: both consumers execute in child processes,// so CPU and peak RSS exclude the fixture server. Latency includes server// scheduling delay; report that delay rather than attributing it to the client.type loadOutput struct { strings.Builder ready atomic.Bool}
func (o *loadOutput) Write(p []byte) (int, error) { if bytes.Contains(p, []byte("READY\n")) { o.ready.Store(true) } return o.Builder.Write(p)}
type loadFixture struct { count, rate int frames [][]byte compressed bool mu sync.Mutex bytes int64 lateNS int64 err error}
func loadPayload(seq int, epoch int64) []byte { return []byte(fmt.Sprintf(`{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"app.bsky.feed.post","did":"did:plc:fixture","operation":"create","record":{"$type":"app.bsky.feed.post","text":%q},"rev":"rev1","rkey":"%d","seq":%d,"time":"2026-07-13T00:00:01.000000Z"}}`, strings.Repeat("a representative synthetic post with repeated words. ", 8), epoch, seq))}func (f *loadFixture) ServeHTTP(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() if r.URL.Path == "/xrpc/network.bsky.jetstream.getZstdDictionary" { _, f.err = w.Write(dictionary) return } if r.URL.Path != "/xrpc/network.bsky.jetstream.subscribeEvents" { http.NotFound(w, r) return } c, err := websocket.Accept(w, r, &websocket.AcceptOptions{Subprotocols: []string{"xrpc.v1.json"}}) if err != nil { f.err = err return } defer c.CloseNow() ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second) defer cancel() enc, err := zstd.NewWriter(nil, zstd.WithEncoderDict(dictionary), zstd.WithEncoderConcurrency(1)) if err != nil { f.err = err return } defer enc.Close() compressed := r.URL.Query().Get("zstdDictionary") != "" if compressed != f.compressed { f.err = fmt.Errorf("unexpected compression mode") return } start := time.Now() for seq := 1; seq <= f.count; seq++ { scheduled := start.Add(time.Duration(seq-1) * time.Second / time.Duration(f.rate)) if wait := time.Until(scheduled); wait > 0 { timer := time.NewTimer(wait) select { case <-timer.C: case <-ctx.Done(): timer.Stop() f.err = ctx.Err() return } } late := time.Since(scheduled).Nanoseconds() if late > f.lateNS { f.lateNS = late } payload := f.frames[seq-1] typ := websocket.MessageText if compressed { typ = websocket.MessageBinary } // Only the first frame carries the epoch; all others are pre-encoded. if seq == 1 { payload = loadPayload(seq, start.UnixNano()) if compressed { payload = enc.EncodeAll(payload, nil) } } if err = c.Write(ctx, typ, payload); err != nil { f.err = err return } f.bytes += int64(len(payload)) } _, _, _ = c.Read(ctx)}func loadClient(args []string) error { if len(args) != 5 { return fmt.Errorf("load-client HOST COUNT MODE DELAY_NS RATE") } count, err := strconv.Atoi(args[1]) if err != nil || count < 1 || count > 1_000_000 { return fmt.Errorf("invalid count") } delay, err := strconv.ParseInt(args[3], 10, 64) if err != nil || delay < 0 { return fmt.Errorf("invalid delay") } rate, err := strconv.Atoi(args[4]) if err != nil || rate <= 0 { return fmt.Errorf("invalid rate") } var epoch int64 samples := make([]int64, 0, count) fmt.Println("READY") client, err := js.Subscribe(args[0], js.WithBatchSize(1), js.WithZstdCompression(args[2] == "compressed")) if err != nil { return err } defer client.Close() ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() for batch, err := range client.Events(ctx) { if err != nil { return err } for _, ev := range batch.Events() { if ev.Seq != uint64(len(samples)+1) || ev.Commit == nil { return fmt.Errorf("out-of-order or non-commit event") } scheduled, err := strconv.ParseInt(ev.Commit.Rkey, 10, 64) if err != nil { return err } if len(samples) == 0 { epoch = scheduled } scheduled = epoch + int64(len(samples))*int64(time.Second)/int64(rate) samples = append(samples, time.Now().UnixNano()-scheduled) if delay > 0 { time.Sleep(time.Duration(delay)) } } if len(samples) >= count { break } } if len(samples) != count { return fmt.Errorf("incomplete delivery: %d/%d", len(samples), count) } sort.Slice(samples, func(i, j int) bool { return samples[i] < samples[j] }) fmt.Printf("LOAD %d %d %d %d\n", count, samples[count/2], samples[min(count-1, count*99/100)], samples[count-1]) return nil}
type loadResult struct { Profile string `json:"profile"` Client string `json:"client"` Mode string `json:"mode"` Repetition int `json:"repetition"` Count int `json:"count"` Rate int `json:"offered_events_per_second"` DelayNS int64 `json:"handler_delay_ns"` WallSeconds float64 `json:"wall_seconds"` CPUSeconds float64 `json:"cpu_seconds"` PeakRSSBytes int64 `json:"peak_rss_bytes"` P50NS int64 `json:"lag_p50_ns"` P99NS int64 `json:"lag_p99_ns"` MaxNS int64 `json:"lag_max_ns"` PayloadBytes int64 `json:"payload_bytes"` ServerLateNS int64 `json:"server_max_schedule_lateness_ns"`}
func measureLoad(client, mode string, repeat, count, rate int, delay int64) (loadResult, error) { f := &loadFixture{count: count, rate: rate, compressed: mode == "compressed"} enc, err := zstd.NewWriter(nil, zstd.WithEncoderDict(dictionary), zstd.WithEncoderConcurrency(1)) if err != nil { return loadResult{}, err } for seq := 1; seq <= count; seq++ { p := loadPayload(seq, 0) if f.compressed { p = enc.EncodeAll(p, nil) } f.frames = append(f.frames, p) } enc.Close() server := httptest.NewServer(f) defer server.Close() ctx, cancel := context.WithTimeout(context.Background(), 65*time.Second) defer cancel() executable, err := os.Executable() if err != nil { return loadResult{}, err } args := []string{"load-client", server.URL, strconv.Itoa(count), mode, strconv.FormatInt(delay, 10), strconv.Itoa(rate)} if client == "zig" { executable = "zig-out/bin/zig-client" args = []string{server.URL, "load", strconv.Itoa(count), mode, strconv.FormatInt(delay, 10), strconv.Itoa(rate)} } measured, err := measureClient(ctx, executable, args) if err != nil { return loadResult{}, err } out := measured.Output var n int var p50, p99, maxLag int64 if _, err = fmt.Sscanf(strings.TrimSpace(strings.TrimPrefix(string(out), "READY\n")), "LOAD %d %d %d %d", &n, &p50, &p99, &maxLag); err != nil || n != count { return loadResult{}, fmt.Errorf("bad child result: %s", out) } f.mu.Lock() defer f.mu.Unlock() if f.err != nil { return loadResult{}, f.err } return loadResult{"", client, mode, repeat, count, rate, delay, measured.WallSeconds, measured.CPUSeconds, measured.PeakRSSBytes, p50, p99, maxLag, f.bytes, f.lateNS}, nil}
type clientProcess struct { Output string WallSeconds, CPUSeconds float64 PeakRSSBytes int64}
func measureClient(ctx context.Context, executable string, args []string) (clientProcess, error) { cmd := exec.CommandContext(ctx, executable, args...) start := time.Now() var output loadOutput cmd.Stdout = &output cmd.Stderr = &output if err := cmd.Start(); err != nil { return clientProcess{}, err } peakDone := make(chan int64, 1) stopPeak := make(chan struct{}) go func() { var peak int64 ticker := time.NewTicker(time.Millisecond) defer ticker.Stop() for { select { case <-stopPeak: peakDone <- peak return case <-ticker.C: if runtime.GOOS == "linux" && output.ready.Load() { exe, _ := os.Readlink(fmt.Sprintf("/proc/%d/exe", cmd.Process.Pid)) if filepath.Base(exe) != filepath.Base(executable) { continue } data, _ := os.ReadFile(fmt.Sprintf("/proc/%d/status", cmd.Process.Pid)) for _, line := range strings.Split(string(data), "\n") { var kb int64 if _, e := fmt.Sscanf(line, "VmHWM: %d kB", &kb); e == nil && kb*1024 > peak { peak = kb * 1024 } } } } } }() err := cmd.Wait() close(stopPeak) observedPeak := <-peakDone out := output.String() wall := time.Since(start).Seconds() if err != nil { return clientProcess{}, fmt.Errorf("client process: %w: %s", err, out) } usage, ok := cmd.ProcessState.SysUsage().(*syscall.Rusage) if !ok { return clientProcess{}, fmt.Errorf("resource usage unavailable") } rss := usage.Maxrss if runtime.GOOS != "darwin" { rss = observedPeak if rss == 0 { return clientProcess{}, fmt.Errorf("no post-exec Linux memory sample") } } return clientProcess{out, wall, cmd.ProcessState.UserTime().Seconds() + cmd.ProcessState.SystemTime().Seconds(), rss}, nil}func loadMain(args []string) error { if len(args) == 0 { return fmt.Errorf("expected load [ci] or catchup [ci|ARCHIVE_COUNT]") } if args[0] == "catchup-client" { return catchupClient(args[1:]) } if args[0] == "catchup" { return catchupMain(args[1:]) } if args[0] == "load-client" { return loadClient(args[1:]) } if args[0] != "load" { return fmt.Errorf("expected load or load-client") } if len(args) > 2 || len(args) == 2 && args[1] != "ci" { return fmt.Errorf("load [ci]") } profile, repetitions := "full", 4 type workload struct { count, rate int delay int64 } workloads := []workload{{4000, 1000, 0}, {20000, 10000, 0}, {1000, 1000, 2_000_000}} if len(args) == 2 { profile, repetitions = "ci", 2 workloads = []workload{{1000, 1000, 0}, {10000, 10000, 0}, {100, 1000, 2_000_000}} } for rep := 0; rep < repetitions; rep++ { clients := []string{"go", "zig"} if rep%2 == 1 { clients = []string{"zig", "go"} } for _, work := range workloads { for _, mode := range []string{"plain", "compressed"} { for _, client := range clients { r, err := measureLoad(client, mode, rep, work.count, work.rate, work.delay) if err != nil { return err } r.Profile = profile if err = json.NewEncoder(os.Stdout).Encode(r); err != nil { return err } fmt.Fprintf(os.Stderr, "PASS load %s %s rate=%d delay=%d repeat=%d\n", client, mode, work.rate, work.delay, rep) } } } } return nil}