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/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) +} 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) + }) +}