This repository has no description
1/// heavily inspired by <https://github.com/bluesky-social/atproto/blob/c7f5a868837d3e9b3289f988fee2267789327b06/packages/tap/README.md>
2
3package tapc
4
5import (
6 "bytes"
7 "context"
8 "encoding/base64"
9 "encoding/json"
10 "fmt"
11 "net/http"
12 "net/url"
13 "time"
14
15 "github.com/bluesky-social/indigo/atproto/syntax"
16 "github.com/gorilla/websocket"
17 "tangled.org/core/log"
18)
19
20type Handler interface {
21 OnEvent(ctx context.Context, evt Event) error
22 OnError(ctx context.Context, err error)
23}
24
25type ConnectHandler interface {
26 OnConnect(ctx context.Context)
27}
28
29type Client struct {
30 Url string
31 AdminPassword string
32 HTTPClient *http.Client
33}
34
35func NewClient(url, adminPassword string) Client {
36 return Client{
37 Url: url,
38 AdminPassword: adminPassword,
39 HTTPClient: &http.Client{},
40 }
41}
42
43func (c *Client) AddRepos(ctx context.Context, dids []syntax.DID) error {
44 if len(dids) == 0 {
45 return nil
46 }
47 body, err := json.Marshal(map[string][]syntax.DID{"dids": dids})
48 if err != nil {
49 return err
50 }
51 req, err := http.NewRequestWithContext(ctx, "POST", c.Url+"/repos/add", bytes.NewReader(body))
52 if err != nil {
53 return err
54 }
55 req.SetBasicAuth("admin", c.AdminPassword)
56 req.Header.Set("Content-Type", "application/json")
57
58 resp, err := c.HTTPClient.Do(req)
59 if err != nil {
60 return err
61 }
62 defer resp.Body.Close()
63 if resp.StatusCode != http.StatusOK {
64 return fmt.Errorf("tap: /repos/add failed with status %d", resp.StatusCode)
65 }
66 return nil
67}
68
69func (c *Client) RemoveRepos(ctx context.Context, dids []syntax.DID) error {
70 body, err := json.Marshal(map[string][]syntax.DID{"dids": dids})
71 if err != nil {
72 return err
73 }
74 req, err := http.NewRequestWithContext(ctx, "POST", c.Url+"/repos/remove", bytes.NewReader(body))
75 if err != nil {
76 return err
77 }
78 req.SetBasicAuth("admin", c.AdminPassword)
79 req.Header.Set("Content-Type", "application/json")
80
81 resp, err := c.HTTPClient.Do(req)
82 if err != nil {
83 return err
84 }
85 defer resp.Body.Close()
86 if resp.StatusCode != http.StatusOK {
87 return fmt.Errorf("tap: /repos/remove failed with status %d", resp.StatusCode)
88 }
89 return nil
90}
91
92func (c *Client) Connect(ctx context.Context, handler Handler) error {
93 l := log.FromContext(ctx)
94
95 u, err := url.Parse(c.Url)
96 if err != nil {
97 return err
98 }
99 if u.Scheme == "https" {
100 u.Scheme = "wss"
101 } else {
102 u.Scheme = "ws"
103 }
104 u.Path = "/channel"
105
106 url := u.String()
107 basicAuth := "Basic " + base64.StdEncoding.EncodeToString([]byte("admin:"+c.AdminPassword))
108
109 var backoff int
110 for {
111 select {
112 case <-ctx.Done():
113 return ctx.Err()
114 default:
115 }
116
117 header := http.Header{
118 "Authorization": []string{basicAuth},
119 }
120 conn, res, err := websocket.DefaultDialer.DialContext(ctx, url, header)
121 if err != nil {
122 if backoff < 12 {
123 backoff++
124 }
125 l.Warn("dialing failed", "url", url, "err", err, "backoff", backoff)
126 time.Sleep(time.Duration(5*backoff) * time.Second)
127
128 continue
129 }
130 backoff = 0
131 l.Info("connected to tap service", "subscription_code", res.StatusCode)
132
133 if ch, ok := handler.(ConnectHandler); ok {
134 ch.OnConnect(ctx)
135 }
136
137 if err = c.handleConnection(ctx, conn, handler); err != nil {
138 l.Warn("tap connection failed", "err", err)
139 }
140 }
141}
142
143func (c *Client) handleConnection(ctx context.Context, conn *websocket.Conn, handler Handler) error {
144 l := log.FromContext(ctx)
145
146 defer func() {
147 conn.Close()
148 l.Warn("closed tap connection")
149 }()
150 l.Info("established tap connection")
151
152 for {
153 select {
154 case <-ctx.Done():
155 return ctx.Err()
156 default:
157 }
158 _, message, err := conn.ReadMessage()
159 if err != nil {
160 return err
161 }
162
163 var ev Event
164 if err := json.Unmarshal(message, &ev); err != nil {
165 handler.OnError(ctx, fmt.Errorf("failed to parse message: %w", err))
166 continue
167 }
168 if err := handler.OnEvent(ctx, ev); err != nil {
169 handler.OnError(ctx, fmt.Errorf("failed to process event %d: %w", ev.ID, err))
170 continue
171 }
172
173 ack := map[string]any{
174 "type": "ack",
175 "id": ev.ID,
176 }
177 if err := conn.WriteJSON(ack); err != nil {
178 l.Warn("failed to send ack", "err", err)
179 continue
180 }
181 }
182}