Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253package spindle
import ( "context" "errors" "fmt" "os"
"github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/rbac" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets")
// cleanup continues after individual steps failfunc (s *Spindle) WipeRepo(ctx context.Context, repoDid syntax.DID, reason string) error { l := s.l.With("repo", repoDid, "reason", reason) l.Warn("wiping repo")
repos, err := s.db.ReposByDid(repoDid) if err != nil { return fmt.Errorf("list repo rows: %w", err) }
var errs []error fail := func(step string, err error) { if err == nil { return } l.Error("wipe step failed", "step", step, "err", err) errs = append(errs, fmt.Errorf("%s: %w", step, err)) }
// cancel work before deleting the rows it references wids, err := s.cancelRepoWorkflows(ctx, repoDid, repos, reason) if err != nil { return fmt.Errorf("cancel workflows: %w", err) }
// collect refs before deleting their rows refs, err := s.db.ListArtifactRefsByRepo(repoDid.String()) fail("list artifact refs", err)
// leave rows as the retry manifest if object cleanup fails objectsFailed := err != nil refSet := map[string]bool{} for _, ref := range refs { refSet[ref] = true for _, err := range s.stores.Delete(ctx, ref) { if err != nil { objectsFailed = true fail("delete artifact "+ref, err) } } } if !objectsFailed { fail("delete artifact refs", s.db.DeleteArtifactRefsByRepo(repoDid.String())) }
leaseIDs, err := s.db.DeleteMillLeasesByRepo(repoDid.String()) fail("delete mill leases", err)
fail("delete queued jobs", s.db.DeleteJobsByRepo(ctx, repoDid.String())) if s.scheduler != nil { fail("delete schedules", s.scheduler.RemoveRepo(ctx, repoDid.String())) } fail("delete quota state", s.db.DeleteQuotaStateForRepo(ctx, repoDid.String())) fail("delete secrets", s.vault.RemoveAllSecrets(ctx, secrets.RepoIdentifier(repoDid.String()))) fail("delete webhooks", s.db.DeleteWebhooksByRepo(repoDid)) fail("delete pull rounds", s.db.DeletePullRoundsByRepo(repoDid))
collabs, err := s.db.ListCollaboratorsByRepoDid(repoDid) fail("list collaborators", err) for _, c := range collabs { fail("remove collaborator acl", s.e.RemoveCollaborator(c.Subject.String(), rbac.ThisServer, repoDid.String())) } fail("delete collaborator rows", s.db.DeleteRepoCollaboratorsByRepoDid(repoDid))
for _, r := range repos { fail("remove repo acl", s.e.RemoveRepo(r.Owner.String(), rbac.ThisServer, repoDid.String())) }
// use the recorded knot and rkey for status events seen := map[db.PipelineKey]bool{} var keys []db.PipelineKey for _, wid := range wids { k := db.PipelineKey{Knot: wid.PipelineId.Knot, Rkey: wid.PipelineId.Rkey} if !seen[k] { seen[k] = true keys = append(keys, k) } } fail("delete events", s.db.DeleteEventsByRepo(repoDid.String(), keys))
// keep the repo row when cleanup needs a retry if len(errs) > 0 { l.Error("keeping repo row for retry: earlier steps failed") } else { fail("delete repo rows", s.db.DeleteReposByDid(repoDid)) }
for _, leaseID := range leaseIDs { ref := "logs/" + leaseID + ".log" if refSet[ref] { continue } for _, err := range s.stores.Delete(ctx, ref) { fail("delete log "+leaseID, err) } }
fail("delete clone", os.RemoveAll(s.newRepoPath(repoDid)))
for _, r := range repos { fail("release owner interest", s.releaseOwnerInterest(ctx, r.Owner)) }
if len(errs) > 0 { return fmt.Errorf("wiped %s with %d step errors: %w", repoDid, len(errs), errors.Join(errs...)) } l.Info("repo wiped", "rows", len(repos)) return nil}
func (s *Spindle) WipeOwner(ctx context.Context, ownerDid syntax.DID, reason string) error { repos, err := s.db.AllRepos() if err != nil { return fmt.Errorf("list repos: %w", err) } var errs []error for _, r := range repos { if r.Owner != ownerDid || r.RepoDid == "" { continue } if err := s.WipeRepo(ctx, r.RepoDid, reason); err != nil { errs = append(errs, err) } } return errors.Join(errs...)}
func (s *Spindle) cancelRepoWorkflows(ctx context.Context, repoDid syntax.DID, repos []db.Repo, reason string) ([]models.WorkflowId, error) { var wids []models.WorkflowId seen := map[models.WorkflowId]bool{} add := func(knot, rkey, name string) { wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: knot, Rkey: rkey}, Name: name} if !seen[wid] { seen[wid] = true wids = append(wids, wid) } }
pipes, err := s.db.ListPipelineWorkflows(repoDid.String()) if err != nil { return nil, fmt.Errorf("list pipeline workflows: %w", err) } for _, p := range pipes { if p.Knot != "" { add(p.Knot, p.Rkey, p.Name) continue } for _, r := range repos { add(r.Knot, p.Rkey, p.Name) } }
leases, err := s.db.ListMillLeases() if err != nil { return nil, fmt.Errorf("list mill leases: %w", err) } for _, lease := range leases { if lease.RepoDID == repoDid.String() { add(lease.Knot, lease.Rkey, lease.Workflow) } }
var errs []error for _, wid := range wids { st, err := s.db.GetStatus(wid) if err == nil && st != nil && models.StatusKind(st.Status).IsFinish() { continue } if err := s.db.StatusCancelled(wid, reason, -1, s.n); err != nil { errs = append(errs, err) } engine.CancelWorkflow(wid) for _, eng := range s.engs { if err := eng.DestroyWorkflow(ctx, wid); err != nil { errs = append(errs, err) } } if err := os.Remove(models.LogFilePath(s.cfg.Server.LogDir, wid)); err != nil && !os.IsNotExist(err) { errs = append(errs, fmt.Errorf("remove live log for %s: %w", wid, err)) } } return wids, errors.Join(errs...)}
func (s *Spindle) releaseOwnerInterest(ctx context.Context, ownerDid syntax.DID) error { if ownerDid == "" { return nil } n, err := s.db.CountReposByOwner(ownerDid) if err != nil { return err } if n > 0 { return nil } members, err := db.CountSpindleMembersBySubject(s.db, ownerDid.String()) if err != nil { return err } if members > 0 { return nil } // the server owner inherits membership via casbin without a members row if ok, err := s.e.IsSpindleMember(ownerDid.String(), rbac.ThisServer); err != nil { return err } else if ok { return nil } if err := db.RemoveDid(s.db, ownerDid.String()); err != nil { return err } s.jc.RemoveDid(ownerDid.String()) if s.tap != nil { return s.tap.RemoveOwnerDIDs(ctx, []syntax.DID{ownerDid}) } return nil}
// retry wipes recorded by bansfunc (s *Spindle) resumeWipes(ctx context.Context) { bans, err := s.db.BanList() if err != nil { s.l.Error("failed to list bans for reconcile", "err", err) return } for _, b := range bans { // try both ban roles err := s.WipeRepo(ctx, b.SubjectDid, "banned") if oerr := s.WipeOwner(ctx, b.SubjectDid, "banned"); oerr != nil { err = errors.Join(err, oerr) } if err != nil { s.l.Error("ban reconcile failed", "subject", b.SubjectDid, "err", err) } }}