From 5d8a24896ac922d65f80ab548b637a3873badc09 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Sat, 13 Jun 2026 10:46:04 +0000 Subject: [PATCH] api/tangled,lexutil: add missing types for DRISL-CBOR event streaming indigo's `lexutil.Client` doesn't support subscription xrpc methods. New `extlexutil.Client` extends the `LexDo` method to support subscription. It won't perform redialing since we don't know which param is cursor and which property is event sequence. since lexgen doesn't generate code subscription xrpc methods, we hand write it. `api/tangled/*_ext.go` files will be treated as non-generated code. Signed-off-by: Seongmin Lee --- .gitattributes | 1 + api/tangled/pipelinesubscribeLogs_ext.go | 98 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ flake.nix | 11 +++++++++-- lexutil/client.go | 241 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ lexutil/lexutil.go | 49 +++++++++++++++++++++++++++++++++++++++++++++++++ 5 file(s) changed, 398 insertion(s)(+), 2 deletion(s)(-) diff --git a/.gitattributes b/.gitattributes --- a/.gitattributes +++ b/.gitattributes @@ -1,2 +1,3 @@ api/tangled/** linguist-generated -diff +api/tangled/*_ext.go -linguist-generated diff flake.lock -diff diff --git a/api/tangled/pipelinesubscribeLogs_ext.go b/api/tangled/pipelinesubscribeLogs_ext.go new file mode 100644 --- /dev/null +++ b/api/tangled/pipelinesubscribeLogs_ext.go @@ -0,0 +1,98 @@ +// extending code generated from sh.tangled.ci.pipeline.subscribeLogs + +package tangled + +import ( + "context" + "fmt" + "io" + + "github.com/bluesky-social/indigo/events" + lexutil "github.com/bluesky-social/indigo/lex/util" + cbg "github.com/whyrusleeping/cbor-gen" + extlexutil "tangled.org/core/lexutil" +) + +// TODO: generate codes below from lexicon +type CiPipelineSubscribeLogs_Event struct { + Error *events.ErrorFrame + Control *CiPipelineSubscribeLogs_Control + Data *CiPipelineSubscribeLogs_Data + + // some private fields for internal routing perf + Preserialized []byte `json:"-" cborgen:"-"` +} + +func (xevt *CiPipelineSubscribeLogs_Event) Serialize(wc io.Writer) error { + header := events.EventHeader{Op: events.EvtKindMessage} + var obj lexutil.CBOR + + switch { + case xevt.Error != nil: + header.Op = events.EvtKindErrorFrame + obj = xevt.Error + case xevt.Control != nil: + header.MsgType = "#control" + obj = xevt.Control + case xevt.Data != nil: + header.MsgType = "#data" + obj = xevt.Data + default: + return fmt.Errorf("unrecognized event kind") + } + + cborWriter := cbg.NewCborWriter(wc) + if err := header.MarshalCBOR(cborWriter); err != nil { + return fmt.Errorf("failed to write header: %w", err) + } + return obj.MarshalCBOR(cborWriter) +} + +func (xevt *CiPipelineSubscribeLogs_Event) Deserialize(r io.Reader) error { + var header events.EventHeader + if err := header.UnmarshalCBOR(r); err != nil { + return fmt.Errorf("reading header: %w", err) + } + switch header.Op { + case events.EvtKindMessage: + switch header.MsgType { + case "#control": + var evt CiPipelineSubscribeLogs_Control + if err := evt.UnmarshalCBOR(r); err != nil { + return fmt.Errorf("reading repoCommit event: %w", err) + } + xevt.Control = &evt + case "#data": + var evt CiPipelineSubscribeLogs_Data + if err := evt.UnmarshalCBOR(r); err != nil { + return fmt.Errorf("reading repoSync event: %w", err) + } + xevt.Data = &evt + } + case events.EvtKindErrorFrame: + var errframe events.ErrorFrame + if err := errframe.UnmarshalCBOR(r); err != nil { + return err + } + xevt.Error = &errframe + default: + return fmt.Errorf("unrecognized event stream type: %d", header.Op) + } + return nil +} + +func CiPipelineSubscribeLogs(ctx context.Context, c extlexutil.LexClient, pipeline string, workflows []string, sched extlexutil.Scheduler[CiPipelineSubscribeLogs_Event]) error { + defer sched.Shutdown() + + params := map[string]any{} + params["pipeline"] = pipeline + params["workflows"] = workflows + + return c.LexDo(ctx, extlexutil.Subscription, "", CiPipelineSubscribeLogsNSID, params, nil, func(ctx context.Context, cr *cbg.CborReader) error { + var evt CiPipelineSubscribeLogs_Event + if err := evt.Deserialize(cr); err != nil { + return err + } + return sched.AddWork(ctx, "", &evt) + }) +} diff --git a/flake.nix b/flake.nix --- a/flake.nix +++ b/flake.nix @@ -515,14 +515,21 @@ rootDir=$(jj --ignore-working-copy root || git rev-parse --show-toplevel) || (echo "error: can't find repo root?"; exit 1) cd "$rootDir" - rm -f api/tangled/* + # *_ext.go are hand-written extensions; never remove or mutate them + find api/tangled -maxdepth 1 -type f -not -name '*_ext.go' -delete lexgen --build-file lexicon-build-config.json lexicons - sed -i.bak 's/\tutil/\/\/\tutil/' api/tangled/* + + # disable type registration temporarily while running cborgen + find api/tangled -maxdepth 1 -name '*.go' -not -name '*_ext.go' -exec \ + sed -i.bak 's/\tutil/\/\/\tutil/' {} + # lexgen generates incomplete Marshaler/Unmarshaler for union types find api/tangled/*.go -not -name "cbor_gen.go" -exec \ sed -i '/^func.*\(MarshalCBOR\|UnmarshalCBOR\)/,/^}/ s/^/\/\/ /' {} + + for f in api/tangled/*_ext.go; do [ -e "''$f" ] && mv "''$f" "''$f.bak"; done ${pkgs.gotools}/bin/goimports -w api/tangled/* CGO_ENABLED=0 go run ./cmd/cborgen/ + for f in api/tangled/*_ext.go.bak; do [ -e "''$f" ] && mv "''$f" "''${f%.bak}"; done + lexgen --build-file lexicon-build-config.json lexicons rm api/tangled/*.bak ''; diff --git a/lexutil/client.go b/lexutil/client.go new file mode 100644 --- /dev/null +++ b/lexutil/client.go @@ -0,0 +1,241 @@ +package lexutil + +import ( + "cmp" + "context" + "fmt" + "log/slog" + "net/http" + "net/url" + "time" + + indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "github.com/carlmjohnson/versioninfo" + "github.com/gorilla/websocket" + cbg "github.com/whyrusleeping/cbor-gen" +) + +const minHealthyConn = 30 * time.Second + +type Client struct { + indigoxrpc.Client + Dialer websocket.Dialer + Logger *slog.Logger +} + +var _ LexClient = (*Client)(nil) + +func makeParams(p map[string]any) url.Values { + params := url.Values{} + for k, v := range p { + if s, ok := v.([]string); ok { + for _, v := range s { + params.Add(k, v) + } + } else { + params.Add(k, fmt.Sprint(v)) + } + } + return params +} + +type processFn func(ctx context.Context, cr *cbg.CborReader) error + +func (c *Client) LexDo(ctx context.Context, method string, inputEncoding string, endpoint string, params map[string]any, bodyData any, out any) error { + switch method { + case Subscription: + if process, ok := out.(processFn); ok { + return c.LexSubscribe(ctx, endpoint, params, process) + } else if redialer, ok := out.(Redialer); ok { + return c.LexSubscribeWithRedialer(ctx, endpoint, params, redialer) + } else { + return fmt.Errorf("unknown output type: %T", out) + } + default: + return c.Client.LexDo(ctx, method, inputEncoding, endpoint, params, bodyData, out) + } +} + +func (c *Client) getHeader() http.Header { + header := http.Header{} + if c.UserAgent != nil { + header.Set("User-Agent", *c.UserAgent) + } else { + header.Set("User-Agent", "extlexutil/"+versioninfo.Short()) + } + if c.Headers != nil { + for k, v := range c.Headers { + header.Set(k, v) + } + } + return header +} + +func (c *Client) LexSubscribe(ctx context.Context, endpoint string, params map[string]any, process func(ctx context.Context, cr *cbg.CborReader) error) error { + logger := cmp.Or(c.Logger, slog.Default().With("system", "events")) + rurl, err := url.Parse(c.Host) + if err != nil { + return err + } + if rurl.Scheme == "http" { + rurl.Scheme = "ws" + } else { + rurl.Scheme = "wss" + } + surl := rurl.JoinPath("/xrpc", endpoint) + surl.RawQuery = makeParams(params).Encode() + + header := c.getHeader() + + u := surl.String() + conn, resp, err := c.Dialer.DialContext(ctx, u, header) + if err != nil { + return fmt.Errorf("%w: %w", ErrDialFailure, err) + } + + logger.Debug("event subscription response", "code", resp.StatusCode, "url", u) + + return c.handleConn(ctx, conn, process) +} + +func (c *Client) LexSubscribeWithRedialer(ctx context.Context, endpoint string, params map[string]any, redialer Redialer) error { + logger := cmp.Or(c.Logger, slog.Default().With("system", "events")) + rurl, err := url.Parse(c.Host) + if err != nil { + return err + } + if rurl.Scheme == "http" { + rurl.Scheme = "ws" + } else { + rurl.Scheme = "wss" + } + surl := rurl.JoinPath("/xrpc", endpoint) + + header := c.getHeader() + + var backoff int + // returns false if the retry budget is exhausted + sleepBackoff := func() bool { + select { + case <-ctx.Done(): + case <-time.After(time.Duration(5+backoff) * time.Second): + } + backoff++ + return backoff <= 15 + } + + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + surl.RawQuery = makeParams(params).Encode() + + u := surl.String() + conn, resp, err := c.Dialer.DialContext(ctx, u, header) + if err != nil { + logger.Warn("dialing failed", "err", err, "backoff", backoff) + if !sleepBackoff() { + return fmt.Errorf("%w: %w", ErrDialFailure, err) + } + continue + } + + logger.Debug("event subscription response", "code", resp.StatusCode, "url", u) + + connectedAt := time.Now() + connErr := c.handleConn(ctx, conn, redialer.Process) + if connErr != nil { + logger.Warn("host connection failed", "err", connErr, "backoff", backoff) + } + + // updates cursor + updated := redialer.UpdateParams(ctx, params) + + // a connection that drops immediately shouldnt reset backoff + // this to avoid reconnect storms + if updated || time.Since(connectedAt) >= minHealthyConn { + backoff = 0 + continue + } + if !sleepBackoff() { + return fmt.Errorf("%w: %w", ErrConnFailure, connErr) + } + } +} + +func (c *Client) handleConn(ctx context.Context, conn *websocket.Conn, process func(ctx context.Context, cr *cbg.CborReader) error) error { + logger := cmp.Or(c.Logger, slog.Default().With("system", "events")) + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + go func() { + t := time.NewTicker(time.Second * 30) + defer t.Stop() + failcount := 0 + + for { + + select { + case <-t.C: + if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second*10)); err != nil { + logger.Warn("failed to ping", "err", err) + failcount++ + if failcount >= 4 { + logger.Error("too many ping fails", "count", failcount) + conn.Close() + return + } + } else { + failcount = 0 // ok ping + } + case <-ctx.Done(): + conn.Close() + return + } + } + }() + + conn.SetPingHandler(func(message string) error { + err := conn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(time.Second*60)) + if err == websocket.ErrCloseSent { + return nil + } + return err + }) + + conn.SetPongHandler(func(_ string) error { + if err := conn.SetReadDeadline(time.Now().Add(time.Minute)); err != nil { + logger.Error("failed to set read deadline", "err", err) + } + + return nil + }) + + cr := new(cbg.CborReader) + + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + mt, rawReader, err := conn.NextReader() + if err != nil { + return fmt.Errorf("conn err at read: %w", err) + } + + if mt != websocket.BinaryMessage { + return fmt.Errorf("expected binary message from subscription endpoint") + } + + cr.SetReader(rawReader) + + if err := process(ctx, cr); err != nil { + return err + } + } +} diff --git a/lexutil/lexutil.go b/lexutil/lexutil.go new file mode 100644 --- /dev/null +++ b/lexutil/lexutil.go @@ -0,0 +1,49 @@ +// extended version of indigo/lex/util package before upstreaming it +package lexutil + +import ( + "context" + "errors" + "io" + + lexutil "github.com/bluesky-social/indigo/lex/util" + cbg "github.com/whyrusleeping/cbor-gen" +) + +type LexClient interface { + lexutil.LexClient + // LexSubscribe is basic event subscriber without redialing logic + // + // golang doesnt allow generics in method so we have to pass raw processFn here instead of Scheduler[T] + LexSubscribe(ctx context.Context, endpoint string, params map[string]any, process func(ctx context.Context, cr *cbg.CborReader) error) error +} + +const Subscription = "subscription" + +var ( + ErrDialFailure = errors.New("dialing failed") + ErrConnFailure = errors.New("connection failed") +) + +type EventStreamMessage interface { + Serialize(wc io.Writer) error + Deserialize(r io.Reader) error +} + +type Scheduler[T any] interface { + AddWork(ctx context.Context, namespace string, val *T) error + Shutdown() +} + +type SeqScheduler[T any] interface { + Scheduler[T] + LastSeq() int64 +} + +type Redialer interface { + // Process decodes the raw message and schedule it + Process(ctx context.Context, cr *cbg.CborReader) error + + // UpdateParams increments the cursor parameter based on LastSeq stored in internal scheduler + UpdateParams(ctx context.Context, params map[string]any) (updated bool) +} -- tangled.sh