This repository has no description
1package state
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "slices"
10 "time"
11
12 "tangled.org/core/appview/cloudflare"
13 "tangled.org/core/appview/notify"
14
15 "tangled.org/core/api/tangled"
16 "tangled.org/core/appview/cache"
17 "tangled.org/core/appview/config"
18 "tangled.org/core/appview/db"
19 "tangled.org/core/appview/models"
20 "tangled.org/core/appview/sites"
21 ec "tangled.org/core/eventconsumer"
22 "tangled.org/core/eventconsumer/cursor"
23 "tangled.org/core/log"
24 "tangled.org/core/orm"
25 "tangled.org/core/rbac"
26 "tangled.org/core/workflow"
27
28 "github.com/bluesky-social/indigo/atproto/syntax"
29 "github.com/go-git/go-git/v5/plumbing"
30 "github.com/posthog/posthog-go"
31)
32
33func Knotstream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, cfClient *cloudflare.Client) (*ec.Consumer, error) {
34 logger := log.FromContext(ctx)
35 logger = log.SubLogger(logger, "knotstream")
36
37 knots, err := db.GetRegistrations(
38 d,
39 orm.FilterIsNot("registered", "null"),
40 )
41 if err != nil {
42 return nil, err
43 }
44
45 srcs := make(map[ec.Source]struct{})
46 for _, k := range knots {
47 s := ec.NewKnotSource(k.Domain)
48 srcs[s] = struct{}{}
49 }
50
51 cache := cache.New(c.Redis.Addr)
52 cursorStore := cursor.NewRedisCursorStore(cache)
53
54 cfg := ec.ConsumerConfig{
55 Sources: srcs,
56 ProcessFunc: knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient),
57 RetryInterval: c.Knotstream.RetryInterval,
58 MaxRetryInterval: c.Knotstream.MaxRetryInterval,
59 ConnectionTimeout: c.Knotstream.ConnectionTimeout,
60 WorkerCount: c.Knotstream.WorkerCount,
61 QueueSize: c.Knotstream.QueueSize,
62 Logger: logger,
63 Dev: c.Core.Dev,
64 CursorStore: &cursorStore,
65 }
66
67 return ec.NewConsumer(cfg), nil
68}
69
70func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) {
71 if repoDid != nil && *repoDid != "" {
72 return db.GetRepoByDid(d, *repoDid)
73 }
74 repos, err := db.GetRepos(d, orm.FilterEq("did", ownerDid), orm.FilterEq("name", repoName))
75 if err != nil {
76 return nil, err
77 }
78 if len(repos) == 0 {
79 return nil, sql.ErrNoRows
80 }
81 return &repos[0], nil
82}
83
84func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client) ec.ProcessFunc {
85 return func(ctx context.Context, source ec.Source, msg ec.Message) error {
86 switch msg.Nsid {
87 case tangled.GitRefUpdateNSID:
88 return ingestRefUpdate(ctx, d, enforcer, posthog, notifier, dev, c, cfClient, source, msg)
89 case tangled.PipelineNSID:
90 return ingestPipeline(d, source, msg)
91 }
92
93 return nil
94 }
95}
96
97func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client, source ec.Source, msg ec.Message) error {
98 logger := log.FromContext(ctx)
99
100 var record tangled.GitRefUpdate
101 err := json.Unmarshal(msg.EventJson, &record)
102 if err != nil {
103 return err
104 }
105
106 knownKnots, err := enforcer.GetKnotsForUser(record.CommitterDid)
107 if err != nil {
108 return err
109 }
110 if !slices.Contains(knownKnots, source.Key()) {
111 return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Key())
112 }
113
114 ownerDid := ""
115 if record.OwnerDid != nil {
116 ownerDid = *record.OwnerDid
117 }
118
119 repo, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName)
120 if lookupErr != nil {
121 return fmt.Errorf("failed to look up repo: %w", lookupErr)
122 }
123
124 logger.Info("processing gitRefUpdate event",
125 "repo", repo.RepoIdentifier(),
126 "ref", record.Ref,
127 "old_sha", record.OldSha,
128 "new_sha", record.NewSha)
129
130 notifier.Push(ctx, repo, record.Ref, record.OldSha, record.NewSha, record.CommitterDid)
131
132 errPunchcard := populatePunchcard(d, record)
133 errLanguages := updateRepoLanguages(d, record)
134
135 var errPosthog error
136 if !dev && record.CommitterDid != "" {
137 errPosthog = pc.Enqueue(posthog.Capture{
138 DistinctId: record.CommitterDid,
139 Event: "git_ref_update",
140 })
141 }
142
143 // Trigger a sites redeploy if this push is to the configured sites branch.
144 if cfClient.Enabled() {
145 go triggerSitesDeployIfNeeded(ctx, d, cfClient, c, record, source)
146 }
147
148 return errors.Join(errPunchcard, errLanguages, errPosthog)
149}
150
151// triggerSitesDeployIfNeeded checks whether the pushed ref matches the sites
152// branch configured for this repo and, if so, syncs the site to R2
153func triggerSitesDeployIfNeeded(ctx context.Context, d *db.DB, cfClient *cloudflare.Client, c *config.Config, record tangled.GitRefUpdate, source ec.Source) {
154 logger := log.FromContext(ctx)
155
156 ref := plumbing.ReferenceName(record.Ref)
157 if !ref.IsBranch() {
158 return
159 }
160 pushedBranch := ref.Short()
161
162 ownerDid := ""
163 if record.OwnerDid != nil {
164 ownerDid = *record.OwnerDid
165 }
166
167 repo, err := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName)
168 if err != nil {
169 return
170 }
171
172 siteConfig, err := db.GetRepoSiteConfig(d, repo.RepoAt().String())
173 if err != nil || siteConfig == nil {
174 return
175 }
176 if siteConfig.Branch != pushedBranch {
177 return
178 }
179
180 scheme := "https"
181 if c.Core.Dev {
182 scheme = "http"
183 }
184 knotHost := fmt.Sprintf("%s://%s", scheme, source.Key())
185
186 deploy := &models.SiteDeploy{
187 RepoAt: repo.RepoAt().String(),
188 Branch: siteConfig.Branch,
189 Dir: siteConfig.Dir,
190 CommitSHA: record.NewSha,
191 Trigger: models.SiteDeployTriggerPush,
192 }
193
194 deployErr := sites.Deploy(ctx, cfClient, knotHost, repo.RepoIdentifier(), record.RepoName, siteConfig.Branch, siteConfig.Dir)
195 if deployErr != nil {
196 logger.Error("sites: R2 sync failed on push", "repo", repo.RepoIdentifier(), "err", deployErr)
197 deploy.Status = models.SiteDeployStatusFailure
198 deploy.Error = deployErr.Error()
199 } else {
200 deploy.Status = models.SiteDeployStatusSuccess
201 }
202
203 if err := db.AddSiteDeploy(d, deploy); err != nil {
204 logger.Error("sites: failed to record deploy", "repo", repo.RepoIdentifier(), "err", err)
205 }
206
207 if deployErr == nil {
208 logger.Info("site deployed to r2", "repo", repo.RepoIdentifier())
209 }
210}
211
212func populatePunchcard(d *db.DB, record tangled.GitRefUpdate) error {
213 if record.CommitterDid == "" {
214 return nil
215 }
216
217 knownEmails, err := db.GetAllEmails(d, record.CommitterDid)
218 if err != nil {
219 return err
220 }
221
222 count := 0
223 for _, ke := range knownEmails {
224 if record.Meta == nil {
225 continue
226 }
227 if record.Meta.CommitCount == nil {
228 continue
229 }
230 for _, ce := range record.Meta.CommitCount.ByEmail {
231 if ce == nil {
232 continue
233 }
234 if ce.Email == ke.Address || ce.Email == record.CommitterDid {
235 count += int(ce.Count)
236 }
237 }
238 }
239
240 punch := models.Punch{
241 Did: record.CommitterDid,
242 Date: time.Now(),
243 Count: count,
244 }
245 return db.AddPunch(d, punch)
246}
247
248func updateRepoLanguages(d *db.DB, record tangled.GitRefUpdate) error {
249 if record.Meta == nil || record.Meta.LangBreakdown == nil || record.Meta.LangBreakdown.Inputs == nil {
250 return fmt.Errorf("empty language data for repo: %v/%s", record.OwnerDid, record.RepoName)
251 }
252
253 ownerDid := ""
254 if record.OwnerDid != nil {
255 ownerDid = *record.OwnerDid
256 }
257
258 r, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName)
259 if lookupErr != nil {
260 return fmt.Errorf("failed to look up repo: %w", lookupErr)
261 }
262 repo := *r
263
264 ref := plumbing.ReferenceName(record.Ref)
265 if !ref.IsBranch() {
266 return fmt.Errorf("%s is not a valid reference name", ref)
267 }
268
269 var langs []models.RepoLanguage
270 for _, l := range record.Meta.LangBreakdown.Inputs {
271 if l == nil {
272 continue
273 }
274
275 langs = append(langs, models.RepoLanguage{
276 RepoAt: repo.RepoAt(),
277 Ref: ref.Short(),
278 IsDefaultRef: record.Meta.IsDefaultRef,
279 Language: l.Lang,
280 Bytes: l.Size,
281 })
282 }
283
284 tx, err := d.Begin()
285 if err != nil {
286 return err
287 }
288 defer tx.Rollback()
289
290 // update appview's cache
291 err = db.UpdateRepoLanguages(tx, repo.RepoAt(), ref.Short(), langs)
292 if err != nil {
293 fmt.Printf("failed; %s\n", err)
294 // non-fatal
295 }
296
297 return tx.Commit()
298}
299
300func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error {
301 var record tangled.Pipeline
302 err := json.Unmarshal(msg.EventJson, &record)
303 if err != nil {
304 return err
305 }
306
307 if record.TriggerMetadata == nil {
308 return fmt.Errorf("empty trigger metadata: nsid %s, rkey %s", msg.Nsid, msg.Rkey)
309 }
310
311 if record.TriggerMetadata.Repo == nil {
312 return fmt.Errorf("empty repo: nsid %s, rkey %s", msg.Nsid, msg.Rkey)
313 }
314
315 repoName := ""
316 if record.TriggerMetadata.Repo.Repo != nil {
317 repoName = *record.TriggerMetadata.Repo.Repo
318 }
319
320 repo, lookupErr := resolveRepo(d, record.TriggerMetadata.Repo.RepoDid, record.TriggerMetadata.Repo.Did, repoName)
321 if lookupErr != nil {
322 return fmt.Errorf("failed to look up repo: %w", lookupErr)
323 }
324 if repo.Spindle == "" {
325 return fmt.Errorf("repo does not have a spindle configured yet: nsid %s, rkey %s", msg.Nsid, msg.Rkey)
326 }
327
328 // trigger info
329 var trigger models.Trigger
330 var sha string
331 trigger.Kind = workflow.TriggerKind(record.TriggerMetadata.Kind)
332 switch trigger.Kind {
333 case workflow.TriggerKindPush:
334 trigger.PushRef = &record.TriggerMetadata.Push.Ref
335 trigger.PushNewSha = &record.TriggerMetadata.Push.NewSha
336 trigger.PushOldSha = &record.TriggerMetadata.Push.OldSha
337 sha = *trigger.PushNewSha
338 case workflow.TriggerKindPullRequest:
339 trigger.PRSourceBranch = &record.TriggerMetadata.PullRequest.SourceBranch
340 trigger.PRTargetBranch = &record.TriggerMetadata.PullRequest.TargetBranch
341 trigger.PRSourceSha = &record.TriggerMetadata.PullRequest.SourceSha
342 trigger.PRAction = &record.TriggerMetadata.PullRequest.Action
343 sha = *trigger.PRSourceSha
344 }
345
346 tx, err := d.Begin()
347 if err != nil {
348 return fmt.Errorf("failed to start txn: %w", err)
349 }
350
351 triggerId, err := db.AddTrigger(tx, trigger)
352 if err != nil {
353 return fmt.Errorf("failed to add trigger entry: %w", err)
354 }
355
356 pipeline := models.Pipeline{
357 Rkey: msg.Rkey,
358 Knot: source.Key(),
359 RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did),
360 RepoName: repoName,
361 RepoDid: repo.RepoDid,
362 TriggerId: int(triggerId),
363 Sha: sha,
364 }
365
366 err = db.AddPipeline(tx, pipeline)
367 if err != nil {
368 return fmt.Errorf("failed to add pipeline: %w", err)
369 }
370
371 err = tx.Commit()
372 if err != nil {
373 return fmt.Errorf("failed to commit txn: %w", err)
374 }
375
376 return nil
377}