From 34096cef660d4d980e7bc00e2ea55f2ad62140d2 Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 06 Jul 2026 07:49:43 +0000 Subject: [PATCH] appview/ingester: ingest label op crud via parking Lewis: May this revision serve well! --- appview/ingester.go | 272 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------------------------------------------------------------------------------------------- appview/ingester_label_test.go | 391 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ appview/ingester_repo.go | 6 +++++- appview/ingester_state_test.go | 53 ++++++++++++++++++++++++++++++++++++----------------- appview/db/label.go | 16 ++++++++++++++++ 5 file(s) changed, 615 insertion(s)(+), 123 deletion(s)(-) diff --git a/appview/ingester.go b/appview/ingester.go --- a/appview/ingester.go +++ b/appview/ingester.go @@ -8,7 +8,6 @@ "fmt" "io" "log/slog" - "maps" "net/http" "net/url" "slices" @@ -1409,7 +1408,7 @@ return err } - i.drainPendingState(ctx, issue.AtUri(), issueStateSpec, l) + i.drainPendingState(ctx, issue.AtUri(), l) l.Info("ingested record") return nil @@ -1576,7 +1575,7 @@ return err } - i.drainPendingState(ctx, pull.AtUri(), pullStatusSpec, l) + i.drainPendingState(ctx, pull.AtUri(), l) l.Info("ingested record") return nil @@ -1815,20 +1814,44 @@ return err } - l.Info("parked state record for retry", "subject", subject) + l.Info("parked record for retry", "subject", subject, "nsid", nsid) return nil } -func (i *Ingester) drainPendingState(ctx context.Context, subject syntax.ATURI, spec stateIngestSpec, l *slog.Logger) { +func (i *Ingester) drainPendingState(ctx context.Context, subject syntax.ATURI, l *slog.Logger) { pending, err := db.PendingStateRecordsForSubject(i.Db, subject) if err != nil { - l.Error("failed to load pending state records", "err", err, "subject", subject) + l.Error("failed to load pending records", "err", err, "subject", subject) return } for _, p := range pending { - if err := i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, spec, l); err != nil { - l.Error("failed to drain pending state record", "err", err, "did", p.Did, "rkey", p.Rkey) + if err := i.reapplyPendingRecord(ctx, p, l); err != nil { + l.Error("failed to drain pending record", "err", err, "did", p.Did, "rkey", p.Rkey, "nsid", p.Nsid) } + } +} + +func (i *Ingester) drainPendingLabelOps(l *slog.Logger) { + subjects, err := db.PendingStateSubjectsForNsid(i.Db, tangled.LabelOpNSID) + if err != nil { + l.Error("failed to list pending label op subjects", "err", err) + return + } + for _, subject := range subjects { + i.drainPendingState(i.Ctx, subject, l) + } +} + +func (i *Ingester) reapplyPendingRecord(ctx context.Context, p db.PendingStateRecord, l *slog.Logger) error { + switch p.Nsid { + case tangled.RepoIssueStateNSID: + return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, issueStateSpec, l) + case tangled.RepoPullStatusNSID: + return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, pullStatusSpec, l) + case tangled.LabelOpNSID: + return i.applyLabelOpRecord(ctx, p.Did, p.Rkey, p.Record, l) + default: + return fmt.Errorf("no reapply handler for parked nsid: %s", p.Nsid) } } @@ -1836,17 +1859,6 @@ pendingStateReconcileInterval = time.Hour pendingStateRecordTTL = 7 * 24 * time.Hour ) - -func stateSpecForSubject(subject syntax.ATURI) (stateIngestSpec, bool) { - switch string(subject.Collection()) { - case tangled.RepoIssueNSID: - return issueStateSpec, true - case tangled.RepoPullNSID: - return pullStatusSpec, true - default: - return stateIngestSpec{}, false - } -} func (i *Ingester) StartPendingStateReconciler() { i.ReconcilePendingState() @@ -1871,11 +1883,7 @@ l.Error("failed to list pending state subjects", "err", err) } for _, subject := range subjects { - spec, ok := stateSpecForSubject(subject) - if !ok { - continue - } - i.drainPendingState(i.Ctx, subject, spec, l) + i.drainPendingState(i.Ctx, subject, l) } cutoff := time.Now().Add(-pendingStateRecordTTL).UTC().Format(time.RFC3339) @@ -2081,6 +2089,10 @@ return fmt.Errorf("failed to create labeldef: %w", err) } + if e.Commit.Operation == jmodels.CommitOperationCreate { + i.drainPendingLabelOps(l) + } + l.Info("ingested record") return nil @@ -2101,92 +2113,142 @@ } func (i *Ingester) ingestLabelOp(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { + l = l.With("handler", "ingestLabelOp") did := e.Did rkey := e.Commit.RKey - var err error - - l = l.With("handler", "ingestLabelOp") - switch e.Commit.Operation { - case jmodels.CommitOperationCreate: - raw := json.RawMessage(e.Commit.Record) - record := tangled.LabelOp{} - err = json.Unmarshal(raw, &record) - if err != nil { - return fmt.Errorf("invalid record: %w", err) - } - - subject := syntax.ATURI(record.Subject) - collection := subject.Collection() - - var repo *models.Repo - switch collection { - case tangled.RepoIssueNSID: - i, err := db.GetIssues(i.Db, orm.FilterEq("at_uri", subject)) - if err != nil || len(i) != 1 { - return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) - } - repo = i[0].Repo - case tangled.RepoPullNSID: - p, err := db.GetPulls(i.Db, orm.FilterEq("at_uri", subject)) - if err != nil || len(p) != 1 { - return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(p)) - } - repo = p[0].Repo - default: - return fmt.Errorf("unsupported label subject: %s", collection) - } - - actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", repo.Labels)) - if err != nil { - return fmt.Errorf("failed to build label application ctx: %w", err) - } - - ops := models.LabelOpsFromRecord(did, rkey, record) - - for _, o := range ops { - def, ok := actx.Defs[o.OperandKey] - if !ok { - return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs))) - } - // validate permissions: only collaborators can apply labels currently - // - // TODO: introduce a repo:triage permission - allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, o.Did, "repo:push") - if permErr != nil { - if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { - return fmt.Errorf("enforcing permission: %w", permErr) - } - l.Warn("ingesting labelop without permission check", "did", o.Did, "err", permErr) - } else if !allowed { - return fmt.Errorf("unauthorized label operation") - } - - if err := def.ValidateOperandValue(&o); err != nil { - return fmt.Errorf("failed to validate labelop: %w", err) - } - } - - tx, err := i.Db.Begin() - if err != nil { - return err - } - defer tx.Rollback() - - for _, o := range ops { - _, err = db.AddLabelOp(tx, &o) - if err != nil { - return fmt.Errorf("failed to add labelop: %w", err) - } - } - - if err = tx.Commit(); err != nil { - return err - } - - l.Info("ingested record") + case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: + return i.applyLabelOpRecord(ctx, did, rkey, e.Commit.Record, l) + case jmodels.CommitOperationDelete: + return i.deleteLabelOpRecord(ctx, did, rkey, l) } + return nil +} + +func (i *Ingester) findLabelSubjectRepo(subject syntax.ATURI) (*models.Repo, bool, error) { + var spec stateIngestSpec + switch subject.Collection() { + case tangled.RepoIssueNSID: + spec = issueStateSpec + case tangled.RepoPullNSID: + spec = pullStatusSpec + default: + return nil, false, fmt.Errorf("unsupported label subject: %s", subject.Collection()) + } + repo, _, found, err := spec.findSubject(i.Db, subject) + return repo, found, err +} + +func (i *Ingester) applyLabelOpRecord(ctx context.Context, did, rkey string, raw []byte, l *slog.Logger) error { + record := tangled.LabelOp{} + if err := json.Unmarshal(raw, &record); err != nil { + return fmt.Errorf("invalid record: %w", err) + } + + subject := syntax.ATURI(record.Subject) + park := func() error { + return i.parkStateRecord(ctx, did, rkey, tangled.LabelOpNSID, subject, raw, l) + } + + repo, found, err := i.findLabelSubjectRepo(subject) + if err != nil { + return err + } + if !found { + return park() + } + + // validate permissions: only collaborators can apply labels currently + // + // TODO: introduce a repo:triage permission + allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") + if permErr != nil { + if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { + return park() + } + l.Warn("ingesting labelop without permission check", "did", did, "err", permErr) + allowed = true + } + + if !allowed { + if err := i.unparkLabelOp(ctx, did, rkey); err != nil { + return err + } + l.Warn("dropped unauthorized label op", "did", did, "rkey", rkey, "subject", subject) + return nil + } + + ops := models.LabelOpsFromRecord(did, rkey, record) + + operandKeys := make([]string, 0, len(ops)) + for idx := range ops { + operandKeys = append(operandKeys, ops[idx].OperandKey) + } + + actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", operandKeys)) + if err != nil { + return fmt.Errorf("failed to build label application ctx: %w", err) + } + + for idx := range ops { + def, ok := actx.Defs[ops[idx].OperandKey] + if !ok { + return park() + } + if err := def.ValidateOperandValue(&ops[idx]); err != nil { + return fmt.Errorf("failed to validate labelop: %w", err) + } + } + + if err := i.materializeLabelOps(ctx, did, rkey, ops); err != nil { + return err + } + + l.Info("ingested record") + return nil +} + +func (i *Ingester) unparkLabelOp(ctx context.Context, did, rkey string) error { + tx, err := i.Db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + + if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { + return fmt.Errorf("failed to unpark label op: %w", err) + } + + return tx.Commit() +} + +func (i *Ingester) materializeLabelOps(ctx context.Context, did, rkey string, ops []models.LabelOp) error { + tx, err := i.Db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + + if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { + return fmt.Errorf("failed to unpark label op: %w", err) + } + if err := db.DeleteLabelOps(tx, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey)); err != nil { + return fmt.Errorf("failed to clear prior label ops: %w", err) + } + for idx := range ops { + if _, err := db.AddLabelOp(tx, &ops[idx]); err != nil { + return fmt.Errorf("failed to add labelop: %w", err) + } + } + return tx.Commit() +} + +func (i *Ingester) deleteLabelOpRecord(ctx context.Context, did, rkey string, l *slog.Logger) error { + if err := i.materializeLabelOps(ctx, did, rkey, nil); err != nil { + return err + } + l.Info("ingested record") return nil } diff --git a/appview/ingester_label_test.go b/appview/ingester_label_test.go new file mode 100644 --- /dev/null +++ b/appview/ingester_label_test.go @@ -0,0 +1,391 @@ +package appview + +import ( + "context" + "encoding/json" + "slices" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + jmodels "github.com/bluesky-social/jetstream/pkg/models" + "tangled.org/core/api/tangled" + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/orm" +) + +const ( + labelOwner = "did:plc:boltless" + labelRepoDid = "did:plc:anemone" + labelPerform = "2026-06-01T00:00:00Z" +) + +func newLabelIngester(t *testing.T) *Ingester { + t.Helper() + ing := newStateIngester(t) + ing.Acl = stubAcl{allow: true} + return ing +} + +func seedLabelDef(t *testing.T, d *db.DB, repoDid, defDid, defRkey, name string, scope []string, multiple, subscribe bool) syntax.ATURI { + t.Helper() + def := &models.LabelDefinition{ + Did: defDid, Rkey: defRkey, Name: name, + ValueType: models.ValueType{Type: models.ConcreteTypeString}, + Scope: scope, Multiple: multiple, Created: time.Now(), + } + if _, err := db.AddLabelDefinition(d, def); err != nil { + t.Fatalf("AddLabelDefinition: %v", err) + } + if subscribe { + if err := db.SubscribeLabel(d, &models.RepoLabel{RepoDid: syntax.DID(repoDid), LabelAt: def.AtUri()}); err != nil { + t.Fatalf("SubscribeLabel: %v", err) + } + } + return def.AtUri() +} + +func labelOpEvent(t *testing.T, op, did, rkey, subject, performedAt string, add, del [][2]string) *jmodels.Event { + t.Helper() + mk := func(pairs [][2]string) []*tangled.LabelOp_Operand { + out := make([]*tangled.LabelOp_Operand, 0, len(pairs)) + for _, p := range pairs { + out = append(out, &tangled.LabelOp_Operand{Key: p[0], Value: p[1]}) + } + return out + } + raw, err := json.Marshal(tangled.LabelOp{Subject: subject, PerformedAt: performedAt, Add: mk(add), Delete: mk(del)}) + if err != nil { + t.Fatalf("marshal: %v", err) + } + return &jmodels.Event{Did: did, Kind: jmodels.EventKindCommit, Commit: &jmodels.Commit{ + Operation: op, Collection: tangled.LabelOpNSID, RKey: rkey, Record: raw}} +} + +func labelDefEvent(t *testing.T, did, rkey, name string, scope []string, multiple bool) *jmodels.Event { + t.Helper() + vt := tangled.LabelDefinition_ValueType{Type: string(models.ConcreteTypeString), Format: string(models.ValueTypeFormatAny)} + raw, err := json.Marshal(tangled.LabelDefinition{Name: name, Scope: scope, Multiple: &multiple, CreatedAt: labelPerform, ValueType: &vt}) + if err != nil { + t.Fatalf("marshal: %v", err) + } + return &jmodels.Event{Did: did, Kind: jmodels.EventKindCommit, Commit: &jmodels.Commit{ + Operation: jmodels.CommitOperationCreate, Collection: tangled.LabelDefinitionNSID, RKey: rkey, Record: raw}} +} + +func mustIngestOp(t *testing.T, ing *Ingester, e *jmodels.Event) { + t.Helper() + if err := ing.ingestLabelOp(context.Background(), e, ing.Logger); err != nil { + t.Fatalf("ingestLabelOp: %v", err) + } +} + +func issueAt(owner, rkey string) syntax.ATURI { + return syntax.ATURI("at://" + owner + "/" + tangled.RepoIssueNSID + "/" + rkey) +} + +func subjectLabels(t *testing.T, d *db.DB, subject syntax.ATURI) []string { + t.Helper() + states, err := db.GetLabels(d, orm.FilterEq("subject", subject)) + if err != nil { + t.Fatalf("GetLabels: %v", err) + } + vals := states[subject].LabelNameValues() + slices.Sort(vals) + return vals +} + +func labelOpValues(t *testing.T, d *db.DB, did, rkey string) []string { + t.Helper() + ops, err := db.GetLabelOps(d, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey)) + if err != nil { + t.Fatalf("GetLabelOps: %v", err) + } + out := make([]string, 0, len(ops)) + for _, op := range ops { + out = append(out, op.OperandValue) + } + slices.Sort(out) + return out +} + +func pendingCount(t *testing.T, d *db.DB, subject syntax.ATURI) int { + t.Helper() + p, err := db.PendingStateRecordsForSubject(d, subject) + if err != nil { + t.Fatalf("PendingStateRecordsForSubject: %v", err) + } + return len(p) +} + +type labelStep struct { + op, rkey, subj, perform string + add, del []string +} + +type foldCase struct { + name string + defRkey string + defName string + scope []string + multiple bool + unsub bool + steps []labelStep + want map[string][]string + reversible bool +} + +func runFold(t *testing.T, tc foldCase, steps []labelStep) map[string][]string { + t.Helper() + ing := newLabelIngester(t) + defUri := "at://" + labelOwner + "/" + tangled.LabelDefinitionNSID + "/" + tc.defRkey + + seedRepo(t, ing.Db, labelOwner, labelRepoDid) + seeded := map[string]bool{} + for _, s := range steps { + if !seeded[s.subj] { + seedIssue(t, ing.Db, labelOwner, labelRepoDid, s.subj) + seeded[s.subj] = true + } + } + seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, tc.defRkey, tc.defName, tc.scope, tc.multiple, !tc.unsub) + + pairs := func(vals []string) [][2]string { + out := make([][2]string, 0, len(vals)) + for _, v := range vals { + out = append(out, [2]string{defUri, v}) + } + return out + } + kind := map[string]string{"create": jmodels.CommitOperationCreate, "update": jmodels.CommitOperationUpdate, "delete": jmodels.CommitOperationDelete} + for _, s := range steps { + perform := s.perform + if perform == "" { + perform = labelPerform + } + mustIngestOp(t, ing, labelOpEvent(t, kind[s.op], labelOwner, s.rkey, string(issueAt(labelOwner, s.subj)), perform, pairs(s.add), pairs(s.del))) + } + + got := map[string][]string{} + for subj := range tc.want { + at := issueAt(labelOwner, subj) + got[subj] = subjectLabels(t, ing.Db, at) + if n := pendingCount(t, ing.Db, at); n != 0 { + t.Fatalf("fold case %q must not park %s, pending=%d", tc.name, subj, n) + } + } + return got +} + +func TestIngestLabelOp_Fold(t *testing.T) { + issue := []string{tangled.RepoIssueNSID} + pull := []string{tangled.RepoPullNSID} + + tests := []foldCase{{ + name: "delete removes the label", defRkey: "prio", defName: "priority", scope: issue, multiple: true, + steps: []labelStep{{op: "create", rkey: "op1", subj: "issue1", add: []string{"high"}}, {op: "delete", rkey: "op1", subj: "issue1"}}, + want: map[string][]string{"issue1": nil}, + }, { + name: "update drops operands the new record no longer carries", defRkey: "prio", defName: "priority", scope: issue, multiple: true, + steps: []labelStep{{op: "create", rkey: "op1", subj: "issue1", add: []string{"high", "low"}}, {op: "update", rkey: "op1", subj: "issue1", add: []string{"high"}}}, + want: map[string][]string{"issue1": {"priority:high"}}, + }, { + name: "update moving subject leaves no stale label", defRkey: "prio", defName: "priority", scope: issue, multiple: true, + steps: []labelStep{{op: "create", rkey: "op1", subj: "issue1", add: []string{"high"}}, {op: "update", rkey: "op1", subj: "issue2", add: []string{"high"}}}, + want: map[string][]string{"issue1": nil, "issue2": {"priority:high"}}, + }, { + name: "a pull-scoped label on an issue is dropped by the fold", defRkey: "prio", defName: "priority", scope: pull, multiple: true, + steps: []labelStep{{op: "create", rkey: "op1", subj: "issue1", add: []string{"high"}}}, + want: map[string][]string{"issue1": nil}, reversible: true, + }, { + name: "a globally-defined but unsubscribed def still applies", defRkey: "prio", defName: "priority", scope: issue, multiple: true, unsub: true, + steps: []labelStep{{op: "create", rkey: "op1", subj: "issue1", add: []string{"high"}}}, + want: map[string][]string{"issue1": {"priority:high"}}, reversible: true, + }, { + name: "a same-instant tie resolves by source rkey", defRkey: "st", defName: "status", scope: issue, multiple: false, + steps: []labelStep{{op: "create", rkey: "aaa", subj: "issue1", add: []string{"open"}}, {op: "create", rkey: "bbb", subj: "issue1", add: []string{"closed"}}}, + want: map[string][]string{"issue1": {"status:closed"}}, reversible: true, + }} + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + orders := [][]labelStep{tc.steps} + if tc.reversible { + rev := slices.Clone(tc.steps) + slices.Reverse(rev) + orders = append(orders, rev) + } + for _, steps := range orders { + got := runFold(t, tc, steps) + for subj, want := range tc.want { + if !slices.Equal(got[subj], want) { + t.Fatalf("subject %s: got %v want %v", subj, got[subj], want) + } + } + } + }) + } +} + +func TestIngestLabelOp_ReingestIsIdempotent(t *testing.T) { + ing := newLabelIngester(t) + subject := seedRepoAndIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + def := seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, "prio", "priority", []string{tangled.RepoIssueNSID}, true, true) + e := labelOpEvent(t, jmodels.CommitOperationCreate, labelOwner, "op1", string(subject), labelPerform, + [][2]string{{def.String(), "high"}, {def.String(), "low"}}, nil) + + mustIngestOp(t, ing, e) + first := labelOpValues(t, ing.Db, labelOwner, "op1") + mustIngestOp(t, ing, e) + if second := labelOpValues(t, ing.Db, labelOwner, "op1"); !slices.Equal(first, second) { + t.Fatalf("re-ingesting an unchanged record changed its stored ops: %v -> %v", first, second) + } + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high", "priority:low"}) { + t.Fatalf("re-ingest changed the derived set: %v", got) + } +} + +func TestIngestLabelOp_Parking(t *testing.T) { + defUri := "at://" + labelOwner + "/" + tangled.LabelDefinitionNSID + "/prio" + issue := []string{tangled.RepoIssueNSID} + seedDef := func(ing *Ingester) syntax.ATURI { + return seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, "prio", "priority", issue, true, true) + } + + tests := []struct { + name string + missing string + acl bool + author string + trigger string + wantLabels []string + }{ + {"missing subject parks, drains on reconcile", "subject", true, labelOwner, "seed-subject", []string{"priority:high"}}, + {"missing def parks, drains on reconcile", "def", true, labelOwner, "seed-def", []string{"priority:high"}}, + {"def create drains a parked op without the sweep", "def", true, labelOwner, "def-event", []string{"priority:high"}}, + {"deleting a parked op unparks it and a later def does not resurrect it", "def", true, labelOwner, "delete-then-def", nil}, + {"a parked op that fails auth on drain is dropped", "subject", false, "did:plc:squid", "seed-subject", nil}, + {"a missing-def park is evicted by TTL though its subject exists", "def", true, labelOwner, "evict", nil}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ing := newLabelIngester(t) + if !tt.acl { + ing.Acl = stubAcl{} + } + subject := issueAt(labelOwner, "issue1") + seedRepo(t, ing.Db, labelOwner, labelRepoDid) + if tt.missing != "subject" { + seedIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + } + if tt.missing != "def" { + seedDef(ing) + } + + mustIngestOp(t, ing, labelOpEvent(t, jmodels.CommitOperationCreate, tt.author, "op1", string(subject), labelPerform, + [][2]string{{defUri, "high"}}, nil)) + if n := pendingCount(t, ing.Db, subject); n != 1 { + t.Fatalf("op must park while its %s is missing, pending=%d", tt.missing, n) + } + + switch tt.trigger { + case "seed-subject": + seedIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + ing.ReconcilePendingState() + case "seed-def": + seedDef(ing) + ing.ReconcilePendingState() + case "def-event": + if err := ing.ingestLabelDefinition(labelDefEvent(t, labelOwner, "prio", "priority", issue, true), ing.Logger); err != nil { + t.Fatalf("ingest def: %v", err) + } + case "delete-then-def": + mustIngestOp(t, ing, labelOpEvent(t, jmodels.CommitOperationDelete, tt.author, "op1", string(subject), "", nil, nil)) + if n := pendingCount(t, ing.Db, subject); n != 0 { + t.Fatalf("deleting a parked op must unpark it, pending=%d", n) + } + seedDef(ing) + ing.ReconcilePendingState() + case "evict": + if n, err := db.EvictStalePendingStateRecords(ing.Db, "2999-01-01T00:00:00Z"); err != nil { + t.Fatalf("evict: %v", err) + } else if n != 1 { + t.Fatalf("a stale missing-def park must be evicted, got %d", n) + } + } + + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, tt.wantLabels) { + t.Fatalf("after %q: got %v want %v", tt.trigger, got, tt.wantLabels) + } + if n := pendingCount(t, ing.Db, subject); n != 0 { + t.Fatalf("no record must remain parked, pending=%d", n) + } + }) + } +} + +func TestIngestLabelOp_DeauthorizedUpdatePreservesPriorOps(t *testing.T) { + ing := newLabelIngester(t) + author := "did:plc:squid" + subject := seedRepoAndIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + def := seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, "prio", "priority", []string{tangled.RepoIssueNSID}, true, true) + + mustIngestOp(t, ing, labelOpEvent(t, jmodels.CommitOperationCreate, author, "op1", string(subject), labelPerform, + [][2]string{{def.String(), "high"}}, nil)) + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high"}) { + t.Fatalf("authorized create: got %v", got) + } + + ing.Acl = stubAcl{} + mustIngestOp(t, ing, labelOpEvent(t, jmodels.CommitOperationUpdate, author, "op1", string(subject), labelPerform, + [][2]string{{def.String(), "low"}}, nil)) + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high"}) { + t.Fatalf("a now-unauthorized update must neither apply nor erase the prior authorized op, got %v", got) + } +} + +func TestIngestLabelOp_DeauthorizedRedeliveryPreservesPriorOps(t *testing.T) { + ing := newLabelIngester(t) + author := "did:plc:squid" + subject := seedRepoAndIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + def := seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, "prio", "priority", []string{tangled.RepoIssueNSID}, true, true) + + create := labelOpEvent(t, jmodels.CommitOperationCreate, author, "op1", string(subject), labelPerform, + [][2]string{{def.String(), "high"}}, nil) + mustIngestOp(t, ing, create) + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high"}) { + t.Fatalf("authorized create: got %v", got) + } + + ing.Acl = stubAcl{} + mustIngestOp(t, ing, create) + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high"}) { + t.Fatalf("a redelivered create must not erase a label applied while authorized, got %v", got) + } +} + +func TestIngestLabelOp_DefDeleteDropsOps(t *testing.T) { + ing := newLabelIngester(t) + subject := seedRepoAndIssue(t, ing.Db, labelOwner, labelRepoDid, "issue1") + def := seedLabelDef(t, ing.Db, labelRepoDid, labelOwner, "prio", "priority", []string{tangled.RepoIssueNSID}, true, true) + + mustIngestOp(t, ing, labelOpEvent(t, jmodels.CommitOperationCreate, labelOwner, "op1", string(subject), labelPerform, + [][2]string{{def.String(), "high"}}, nil)) + if got := subjectLabels(t, ing.Db, subject); !slices.Equal(got, []string{"priority:high"}) { + t.Fatalf("after create: got %v", got) + } + + del := &jmodels.Event{Did: labelOwner, Kind: jmodels.EventKindCommit, Commit: &jmodels.Commit{ + Operation: jmodels.CommitOperationDelete, Collection: tangled.LabelDefinitionNSID, RKey: "prio"}} + if err := ing.ingestLabelDefinition(del, ing.Logger); err != nil { + t.Fatalf("ingest def delete: %v", err) + } + if rows := labelOpValues(t, ing.Db, labelOwner, "op1"); len(rows) != 0 { + t.Fatalf("deleting a def must cascade-delete its label_ops rows, survived: %v", rows) + } + if got := subjectLabels(t, ing.Db, subject); len(got) != 0 { + t.Fatalf("deleting a def must drop its labels, got %v", got) + } +} diff --git a/appview/ingester_repo.go b/appview/ingester_repo.go --- a/appview/ingester_repo.go +++ b/appview/ingester_repo.go @@ -250,7 +250,11 @@ if err := applyRepoMetadata(tx, current, desired); err != nil { return fmt.Errorf("failed to apply repo metadata: %w", err) } - return tx.Commit() + if err := tx.Commit(); err != nil { + return err + } + + return nil } func (i *Ingester) ingestRepoDelete(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { diff --git a/appview/ingester_state_test.go b/appview/ingester_state_test.go --- a/appview/ingester_state_test.go +++ b/appview/ingester_state_test.go @@ -2,6 +2,7 @@ import ( "context" + "database/sql" "encoding/json" "errors" "log/slog" @@ -41,21 +42,36 @@ return s.allow, s.err } -func seedRepoAndIssue(t *testing.T, d *db.DB, ownerDid, repoDid, issueRkey string) syntax.ATURI { +func seedTx(t *testing.T, d *db.DB, fn func(tx *sql.Tx) error) { t.Helper() tx, err := d.Begin() if err != nil { t.Fatalf("Begin: %v", err) } - if err := db.AddRepo(tx, &models.Repo{ - Did: ownerDid, - Name: "anemone", - Knot: "knot.example", - Rkey: "anemone", - RepoDid: repoDid, - }); err != nil { - t.Fatalf("AddRepo: %v", err) + defer tx.Rollback() + if err := fn(tx); err != nil { + t.Fatalf("seed: %v", err) } + if err := tx.Commit(); err != nil { + t.Fatalf("Commit: %v", err) + } +} + +func seedRepo(t *testing.T, d *db.DB, ownerDid, repoDid string) { + t.Helper() + seedTx(t, d, func(tx *sql.Tx) error { + return db.AddRepo(tx, &models.Repo{ + Did: ownerDid, + Name: "anemone", + Knot: "knot.example", + Rkey: "anemone", + RepoDid: repoDid, + }) + }) +} + +func seedIssue(t *testing.T, d *db.DB, ownerDid, repoDid, issueRkey string) syntax.ATURI { + t.Helper() issue := &models.Issue{ Did: ownerDid, Rkey: issueRkey, @@ -64,13 +80,16 @@ Body: "body", Open: true, } - if err := db.PutIssue(tx, issue); err != nil { - t.Fatalf("PutIssue: %v", err) - } - if err := tx.Commit(); err != nil { - t.Fatalf("Commit: %v", err) - } + seedTx(t, d, func(tx *sql.Tx) error { + return db.PutIssue(tx, issue) + }) return issue.AtUri() +} + +func seedRepoAndIssue(t *testing.T, d *db.DB, ownerDid, repoDid, issueRkey string) syntax.ATURI { + t.Helper() + seedRepo(t, d, ownerDid, repoDid) + return seedIssue(t, d, ownerDid, repoDid, issueRkey) } func issueStateEvent(t *testing.T, op, did, rkey, subject, state, createdAt string) *jmodels.Event { @@ -194,7 +213,7 @@ } at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") - ing.drainPendingState(ctx, at, issueStateSpec, ing.Logger) + ing.drainPendingState(ctx, at, ing.Logger) if !ingestedIssueOpen(t, ing.Db, at) { t.Fatal("a deleted parked record must not apply after the subject arrives") } @@ -216,7 +235,7 @@ } at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") - ing.drainPendingState(ctx, at, issueStateSpec, ing.Logger) + ing.drainPendingState(ctx, at, ing.Logger) if !ingestedIssueOpen(t, ing.Db, at) { t.Fatal("a parked record that fails authorization on drain must not apply") diff --git a/appview/db/label.go b/appview/db/label.go --- a/appview/db/label.go +++ b/appview/db/label.go @@ -223,6 +223,22 @@ return id, nil } +func DeleteLabelOps(e Execer, filters ...orm.Filter) error { + var conditions []string + var args []any + for _, filter := range filters { + conditions = append(conditions, filter.Condition()) + args = append(args, filter.Arg()...) + } + whereClause := "" + if conditions != nil { + whereClause = " where " + strings.Join(conditions, " and ") + } + query := fmt.Sprintf(`delete from label_ops %s`, whereClause) + _, err := e.Exec(query, args...) + return err +} + func GetLabelOps(e Execer, filters ...orm.Filter) ([]models.LabelOp, error) { var labelOps []models.LabelOp var conditions []string -- tangled.sh