This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

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