diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b493271..d3b16886 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,13 @@ ## Unreleased +## 0.15.0 - 2026-10-01 + +**Highlights:** Exact owner-directed archive removal with durable collection exclusions. + +- Add `purge-threads` with exact owner-selected plans, atomic local removal, and durable collection exclusions. Thanks @hannesrudolph. +- Update CrawlKit to v0.16.7. Thanks @vincentkoc. + ## 0.14.0 - 2026-09-30 **Highlights:** New `gitcrawl analytics` commands for publication repair, actor evidence, and continuous collection with durable review-state recovery. diff --git a/docs/governance.md b/docs/governance.md index 44d98463..dd3bd43e 100644 --- a/docs/governance.md +++ b/docs/governance.md @@ -137,3 +137,28 @@ The thread stays open on GitHub; only your local triage view hides it. - It does not edit, label, comment on, or close GitHub issues. Use `gh` for that. - It does not retrain embeddings or reshape the underlying graph — it overlays decisions on top of the algorithm output. - It does not propagate to other gitcrawl installations unless you publish your database via a [portable store](/portable-stores/). + +## Permanent owner removal + +For a native local archive, preview an exact selection, then apply its `plan_id`: + +```sh +gitcrawl purge-threads owner/repo --numbers 123,456 --json +gitcrawl purge-threads owner/repo --numbers 123,456 --apply PLAN_ID --json +``` + +The read-only plan lists thread identities and affected row counts. Apply rechecks +its content hash in one transaction, removes selected bodies, history, derived +records and retry work, and installs permanent collection exclusions. Changed +plans require a new preview. Shared actor profiles and unrelated threads stay +intact. Owner exclusions are reported separately from provider deletion or +successful recovery; core collection coverage remains independent. + +Selections must contain 1–100 distinct existing numbers in one unambiguous local +repository. The analytics collector must be stopped. Portable/cloud archives, +blob-backed targets, shared cluster state, retained repository workflow snapshots, +linked workflow reservations, and deletions allowing native integer ID reuse are +refused. Exclusions follow the explicit `owner/repo` collection namespace; carry +that policy when changing targets. Portable publication is refused while local +exclusions exist. This removes logical archive data, not GitHub content, backups, +or forensic disk remnants. There is no undo or automatic backup. diff --git a/internal/cli/analytics_integrity.go b/internal/cli/analytics_integrity.go index 03a7adaf..c36c6d93 100644 --- a/internal/cli/analytics_integrity.go +++ b/internal/cli/analytics_integrity.go @@ -379,6 +379,11 @@ func analyticsReviewBudget(limits []gh.RateLimitSnapshot, now time.Time) (int, g } func (a *App) analyticsNumbers(ctx context.Context, s *store.Store, owner, repo string, numbers []int, discovery bool, operation string) (resultErr error) { + var err error + numbers, err = s.FilterExcludedNumbers(ctx, owner+"/"+repo, numbers) + if err != nil { + return err + } var wg sync.WaitGroup var failures []error // Every return, including cancelled admission, waits for durable receipts. diff --git a/internal/cli/app.go b/internal/cli/app.go index 0886bfa3..89cbefea 100644 --- a/internal/cli/app.go +++ b/internal/cli/app.go @@ -141,6 +141,8 @@ func (a *App) Run(ctx context.Context, args []string) error { return a.runThreads(ctx, rest[1:]) case "capture": return a.runCapture(ctx, rest[1:]) + case "purge-threads": + return a.runPurgeThreads(ctx, rest[1:]) case "close-thread": return a.runCloseThread(ctx, rest[1:]) case "reopen-thread": diff --git a/internal/cli/app_test.go b/internal/cli/app_test.go index f2f9ba39..467cde2b 100644 --- a/internal/cli/app_test.go +++ b/internal/cli/app_test.go @@ -4344,10 +4344,10 @@ func TestDoctorJSONReportsCurrentSchemaDiagnosticsWithoutMutation(t *testing.T) if got := schema["state"]; got != "current" { t.Fatalf("db_schema.state = %#v, payload=%#v", got, schema) } - if got := schema["current_version"]; got != float64(15) { + if got := schema["current_version"]; got != float64(16) { t.Fatalf("db_schema.current_version = %#v, payload=%#v", got, schema) } - if got := schema["supported_version"]; got != float64(15) { + if got := schema["supported_version"]; got != float64(16) { t.Fatalf("db_schema.supported_version = %#v, payload=%#v", got, schema) } if got := schema["child_observation_reservations"]; got != true { @@ -4609,7 +4609,7 @@ func TestDoctorJSONReportsLegacyPendingSchemaWithoutMutation(t *testing.T) { t.Fatalf("pr_details.duplicate_path_files_supported = %#v, payload=%#v", got, prDetails) } pending := doctorStringList(t, schema, "pending_migrations") - if !doctorListContains(pending, "schema_version_3_to_15") || + if !doctorListContains(pending, "schema_version_3_to_16") || !doctorListContains(pending, "pull_request_files_position_key") || !doctorListContains(pending, "thread_child_observation_reservations_table") { t.Fatalf("pending_migrations = %#v", pending) diff --git a/internal/cli/gh_search_test.go b/internal/cli/gh_search_test.go index 174a0711..720156db 100644 --- a/internal/cli/gh_search_test.go +++ b/internal/cli/gh_search_test.go @@ -291,8 +291,8 @@ func TestGHSearchSyncIfStaleMigratesFreshPortableRuntime(t *testing.T) { if err := rt.Store.DB().QueryRowContext(ctx, `pragma user_version`).Scan(&schemaVersion); err != nil { t.Fatalf("read runtime schema version: %v", err) } - if schemaVersion != 15 { - t.Fatalf("runtime schema version = %d, want 15", schemaVersion) + if schemaVersion != 16 { + t.Fatalf("runtime schema version = %d, want 16", schemaVersion) } var tableName string if err := rt.Store.DB().QueryRowContext(ctx, `select name from sqlite_schema where type = 'table' and name = 'sync_runs'`).Scan(&tableName); err != nil { diff --git a/internal/cli/help.go b/internal/cli/help.go index 8474dd4f..a0aee183 100644 --- a/internal/cli/help.go +++ b/internal/cli/help.go @@ -65,6 +65,7 @@ Core commands: capture export a stable code-free conversation snapshot code index index tracked text files from a local Git checkout cluster build durable clusters from local thread vectors + purge-threads permanently remove selected local threads and exclude collection close-thread locally hide one issue or pull request row reopen-thread clear a local hide for one issue or pull request row close-cluster locally hide one durable cluster @@ -264,6 +265,17 @@ Usage: Usage: gitcrawl runs owner/repo [--kind sync|summary|embedding|cluster] [--limit N] [--json] `, + "purge-threads": `gitcrawl purge-threads plans owner-directed local removal, never a GitHub deletion. + +Usage: + gitcrawl purge-threads owner/repo --numbers 1,2 [--apply PLAN_ID] [--json] + +Preview the exact selection, then pass its plan_id to --apply. Changed plans +are refused. Requires a native local archive and an idle analytics collector. +Blob-backed targets, workflow dependencies, shared clusters and reusable native +IDs are refused. Exclusions are permanent; portable publication is unavailable. +`, + "close-thread": `gitcrawl close-thread locally hides one issue or pull request row. Usage: diff --git a/internal/cli/thread_purge.go b/internal/cli/thread_purge.go new file mode 100644 index 00000000..598da719 --- /dev/null +++ b/internal/cli/thread_purge.go @@ -0,0 +1,96 @@ +package cli + +import ( + "context" + "flag" + "fmt" + "io" + "os" + "path/filepath" + + "github.com/openclaw/gitcrawl/internal/config" + "github.com/openclaw/gitcrawl/internal/store" +) + +func (a *App) runPurgeThreads(ctx context.Context, args []string) error { + fs := flag.NewFlagSet("purge-threads", flag.ContinueOnError) + fs.SetOutput(io.Discard) + raw := fs.String("numbers", "", "explicit issue/PR numbers") + apply := fs.String("apply", "", "apply the exact previewed plan ID") + jsonOut := fs.Bool("json", false, "JSON plan/result") + if err := fs.Parse(normalizeCommandArgs(args, map[string]bool{"numbers": true, "apply": true})); err != nil { + return usageErr(err) + } + if fs.NArg() != 1 { + return usageErr(fmt.Errorf("purge-threads requires owner/repo")) + } + owner, repo, err := parseOwnerRepo(fs.Arg(0)) + if err != nil { + return usageErr(err) + } + repository := owner + "/" + repo + numbers, err := parseOptionalThreadNumberList(*raw, repository) + if err != nil { + return usageErr(err) + } + if len(numbers) == 0 || len(numbers) > 100 { + return usageErr(fmt.Errorf("purge-threads requires 1..100 explicit --numbers")) + } + if flagWasSet(fs, "apply") && len(*apply) != 64 { + return usageErr(fmt.Errorf("--apply requires the plan_id from a preview")) + } + a.applyCommandJSON(*jsonOut) + cfg, err := config.LoadRuntime(a.configPath) + if err != nil { + return err + } + if cfg.Remote.Enabled() { + return fmt.Errorf("purge-threads requires a native local archive") + } + _, portable, err := portableStoreRoot(ctx, cfg.DBPath) + if err != nil { + return err + } + if portable { + return fmt.Errorf("purge-threads refuses portable stores and mirrors") + } + // Lock before opening a writer: an analytics collector owns this same file. + if *apply != "" { + lock, err := os.OpenFile(filepath.Join(filepath.Dir(cfg.DBPath), "runner.lock"), os.O_CREATE|os.O_RDWR, 0600) + if err != nil { + return err + } + defer lock.Close() + if err = lockPortableFile(lock); err != nil { + return fmt.Errorf("collector ownership lock busy: %w", err) + } + } + probe, err := store.OpenReadOnly(ctx, cfg.DBPath) + if err != nil { + return err + } + plan, err := probe.PlanThreadPurge(ctx, repository, numbers) + closeErr := probe.Close() + if err != nil { + return err + } + if closeErr != nil { + return closeErr + } + if *apply == "" { + return a.writeOutput("thread_purge_plan", plan, true) + } + if *apply != plan.PlanID { + return fmt.Errorf("purge plan changed; preview again before applying") + } + st, err := store.Open(ctx, cfg.DBPath) + if err != nil { + return err + } + defer st.Close() + result, err := st.PurgeThreads(ctx, repository, numbers, *apply) + if err != nil { + return err + } + return a.writeOutput("thread_purge", result, true) +} diff --git a/internal/cli/thread_purge_test.go b/internal/cli/thread_purge_test.go new file mode 100644 index 00000000..abce856f --- /dev/null +++ b/internal/cli/thread_purge_test.go @@ -0,0 +1,93 @@ +package cli + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/openclaw/gitcrawl/internal/store" +) + +func TestThreadPurgeCLIPlansThenRequiresIdleOwner(t *testing.T) { + ctx := context.Background() + dir := t.TempDir() + db := filepath.Join(dir, "archive.db") + cfg := filepath.Join(dir, "config.toml") + if err := os.WriteFile(cfg, []byte(fmt.Sprintf("db_path=%q\n", db)), 0600); err != nil { + t.Fatal(err) + } + s, err := store.Open(ctx, db) + if err != nil { + t.Fatal(err) + } + _, err = s.DB().Exec(`INSERT INTO repositories(id,owner,name,full_name,raw_json,updated_at) VALUES(1,'fixture','repo','fixture/repo','{}','2026-01-01'); + INSERT INTO threads(id,repo_id,github_id,number,kind,state,title,html_url,labels_json,assignees_json,raw_json,content_hash,updated_at) VALUES + (1,1,'one',10,'pull_request','open','PRIVATE_SENTINEL','','[]','[]','{}','h','2026-01-01'), + (2,1,'two',20,'pull_request','open','keeper','','[]','[]','{}','h','2026-01-01')`) + if err != nil { + t.Fatal(err) + } + s.Close() + run := func(args ...string) (string, error) { + a := New() + var out, stderr bytes.Buffer + a.Stdout = &out + a.Stderr = &stderr + err := a.Run(ctx, append([]string{"--config", cfg, "purge-threads", "fixture/repo", "--numbers", "10", "--json"}, args...)) + return out.String(), err + } + out, err := run() + if err != nil || !strings.Contains(out, `"applied": false`) || strings.Contains(out, "PRIVATE_SENTINEL") { + t.Fatalf("plan=%s err=%v", out, err) + } + var plan store.ThreadPurgePlan + if err := json.Unmarshal([]byte(out), &plan); err != nil { + t.Fatal(err) + } + if _, err = run("--apply"); err == nil { + t.Fatal("missing owner request accepted") + } + for _, args := range [][]string{{"--apply", ""}, {"--apply", strings.Repeat("0", 64)}, {"--numbers", ""}, {"--numbers", "10,10"}, {"--numbers", "other/repo#10"}, {"--numbers", "10,"}, {"--numbers", "999"}} { + if _, err := run(args...); err == nil { + t.Fatal("unsafe arguments accepted", args) + } + } + lock, err := os.OpenFile(filepath.Join(dir, "runner.lock"), os.O_CREATE|os.O_RDWR, 0600) + if err != nil { + t.Fatal(err) + } + if err = lockPortableFile(lock); err != nil { + t.Fatal(err) + } + _, err = run("--apply", plan.PlanID) + if err == nil || !strings.Contains(err.Error(), "ownership lock busy") { + t.Fatalf("owner bypass: %v", err) + } + lock.Close() + out, err = run("--apply", plan.PlanID) + if err != nil || !strings.Contains(out, `"applied": true`) { + t.Fatalf("apply=%s %v", out, err) + } + out, err = run() + if err == nil { + t.Fatalf("repeat plan=%s %v", out, err) + } +} + +func TestThreadPurgeRefusesRemoteWithoutAccess(t *testing.T) { + cfg := filepath.Join(t.TempDir(), "config.toml") + if err := os.WriteFile(cfg, []byte(fmt.Sprintf("db_path=%q\n[remote]\nmode=\"cloud\"\nendpoint=\"https://example.invalid\"\narchive=\"fixture\"\n", filepath.Join(filepath.Dir(cfg), "synthetic.db"))), 0600); err != nil { + t.Fatal(err) + } + a := New() + a.Stdout = &bytes.Buffer{} + a.Stderr = &bytes.Buffer{} + if err := a.Run(context.Background(), []string{"--config", cfg, "purge-threads", "fixture/repo", "--numbers", "1"}); err == nil || !strings.Contains(err.Error(), "native local") { + t.Fatal("remote target admitted", err) + } +} diff --git a/internal/store/analytics_integrity.go b/internal/store/analytics_integrity.go index f4046b9e..9418d23c 100644 --- a/internal/store/analytics_integrity.go +++ b/internal/store/analytics_integrity.go @@ -57,6 +57,13 @@ func (s *Store) RecordAnalyticsAttempt(ctx context.Context, a AnalyticsAttempt) return err } return s.WithTx(ctx, func(tx *Store) error { + excluded, err := tx.ThreadExcluded(ctx, a.Repository, a.Number) + if err != nil { + return err + } + if excluded { + return nil + } status, class, message := a.Status, a.ErrorClass, a.ErrorText knownReview := true if status == "success" && (a.Operation == "graphql_history" || a.Operation == "review_state") { @@ -359,12 +366,23 @@ func (s *Store) SaveReviewStateCoverage(ctx context.Context, repository string, if err := s.q().QueryRowContext(ctx, `SELECT count(*) FROM analytics_retries WHERE repository=? AND operation='review_state' AND resolved_at IS NULL`, repository).Scan(&pending); err != nil { return err } - _, err := s.q().ExecContext(ctx, `INSERT INTO analytics_review_state_coverage VALUES(?,?,?,?,?,?,?,?,?) ON CONFLICT(repository) DO UPDATE SET cursor=excluded.cursor,ceiling=excluded.ceiling,scanned=excluded.scanned,queued=excluded.queued,pending_items=excluded.pending_items,scan_complete=excluded.scan_complete,complete=excluded.complete,observed_at=excluded.observed_at`, repository, progress.Cursor, progress.Ceiling, progress.Scanned, progress.Queued, pending, boolInt(progress.Done), boolInt(progress.Done && pending == 0), time.Now().UTC().Format(time.RFC3339Nano)) + excluded, err := s.OwnerExcludedCount(ctx, repository, true) + if err != nil { + return err + } + _, err = s.q().ExecContext(ctx, `INSERT INTO analytics_review_state_coverage(repository,cursor,ceiling,scanned,queued,pending_items,scan_complete,complete,observed_at,owner_excluded_items) VALUES(?,?,?,?,?,?,?,?,?,?) ON CONFLICT(repository) DO UPDATE SET cursor=excluded.cursor,ceiling=excluded.ceiling,scanned=excluded.scanned,queued=excluded.queued,pending_items=excluded.pending_items,scan_complete=excluded.scan_complete,complete=excluded.complete,observed_at=excluded.observed_at,owner_excluded_items=excluded.owner_excluded_items`, repository, progress.Cursor, progress.Ceiling, progress.Scanned, progress.Queued, pending, boolInt(progress.Done), boolInt(progress.Done && pending == 0 && excluded == 0), time.Now().UTC().Format(time.RFC3339Nano), excluded) return err } func (s *Store) AnalyticsIntegrityStatus(ctx context.Context, repository string) (map[string]any, error) { out := map[string]any{"repository": repository, "checked_at": time.Now().UTC().Format(time.RFC3339Nano)} + if s.hasTable(ctx, "thread_exclusions") { + excluded, err := s.OwnerExcludedCount(ctx, repository, false) + if err != nil { + return nil, err + } + out["owner_excluded_items"] = excluded + } var through, observed string var issues, prs, complete int coverageErr := s.q().QueryRowContext(ctx, "SELECT through,issues,pull_requests,complete,observed_at FROM analytics_coverage WHERE repository=?", repository).Scan(&through, &issues, &prs, &complete, &observed) diff --git a/internal/store/analytics_integrity_test.go b/internal/store/analytics_integrity_test.go index 4b16c63d..19d440ab 100644 --- a/internal/store/analytics_integrity_test.go +++ b/internal/store/analytics_integrity_test.go @@ -152,7 +152,7 @@ func TestAnalyticsV14MigrationPreservesReceiptsAndCoverageWatermark(t *testing.T s.DB().QueryRow("PRAGMA user_version").Scan(&version) s.DB().QueryRow("SELECT through,complete FROM analytics_coverage WHERE repository='fixture/repo'").Scan(&through, &complete) cp, err := s.AnalyticsState(ctx, "updates:fixture/repo") - if err != nil || version != 15 || through != "2026-01-01T00:00:00Z" || complete != 1 || cp != `{"kind":1,"cursor":"provider-cursor"}` { + if err != nil || version != schemaVersion || through != "2026-01-01T00:00:00Z" || complete != 1 || cp != `{"kind":1,"cursor":"provider-cursor"}` { t.Fatalf("migration changed evidence: %d %s %d %s %v", version, through, complete, cp, err) } } diff --git a/internal/store/analytics_source.go b/internal/store/analytics_source.go index bb882569..0125c91d 100644 --- a/internal/store/analytics_source.go +++ b/internal/store/analytics_source.go @@ -93,6 +93,11 @@ func (s *Store) queueAnalyticsActor(ctx context.Context, raw string) error { if id == "" { return nil } + if excluded, err := s.NodeExcluded(ctx, id); err != nil { + return err + } else if excluded { + return nil + } _, err := s.q().ExecContext(ctx, "INSERT OR IGNORE INTO analytics_pending_nodes(node_id,kind) VALUES(?,?)", id, kind) return err } @@ -236,6 +241,11 @@ func (s *Store) SaveActorEvidence(ctx context.Context, nodes []map[string]any, a if id == "" { continue } + if excluded, err := tx.NodeExcluded(ctx, id); err != nil { + return err + } else if excluded { + continue + } a, _ := n["author"].(map[string]any) actor, _ := a["id"].(string) if actor != "" { diff --git a/internal/store/portable.go b/internal/store/portable.go index 9fcc4cc9..08dbf0d3 100644 --- a/internal/store/portable.go +++ b/internal/store/portable.go @@ -94,6 +94,17 @@ type PortablePruneStats struct { } func (s *Store) PrunePortablePayloads(ctx context.Context, options PortablePruneOptions) (PortablePruneStats, error) { + // Owner exclusions belong to this local archive. Do not silently drop their + // policy or publish owner-request metadata through a portable export. + if s.hasTable(ctx, "thread_exclusions") { + var n int + if err := s.q().QueryRowContext(ctx, "SELECT count(*) FROM thread_exclusions").Scan(&n); err != nil { + return PortablePruneStats{}, err + } + if n > 0 { + return PortablePruneStats{}, fmt.Errorf("portable publication is unavailable for an archive with local owner exclusions") + } + } if options.BodyChars <= 0 { options.BodyChars = 256 } diff --git a/internal/store/schema_convergence.go b/internal/store/schema_convergence.go index 59698397..7a76520e 100644 --- a/internal/store/schema_convergence.go +++ b/internal/store/schema_convergence.go @@ -150,6 +150,9 @@ func inspectCompatibilityMigrationsMode( if current > 0 && !st.observationSchemaConvergenceHasCurrentShape(ctx) { add(migrationObservationSchemaConvergence) } + if current >= 16 && (!st.hasTable(ctx, "thread_exclusions") || !st.hasTable(ctx, "thread_excluded_nodes") || !st.hasColumn(ctx, "analytics_review_state_coverage", "owner_excluded_items")) { + add("thread_exclusions_schema") + } if !includeSemantic { return pending, nil } diff --git a/internal/store/store.go b/internal/store/store.go index ca464e80..75d2d602 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -17,7 +17,7 @@ import ( ) const ( - schemaVersion = 15 + schemaVersion = 16 timeLayout = time.RFC3339Nano ) @@ -285,9 +285,9 @@ func (s *Store) migrate(ctx context.Context) error { if _, err := s.db.ExecContext(ctx, schemaSQL); err != nil { return fmt.Errorf("apply schema: %w", err) } - // Version 15 only adds operational receipts and an explicit review-thread - // membership column. A converged v14 archive needs no row/history rebuild. - if current == 14 { + // Additive analytics/policy upgrades need no row/history rebuild on a + // converged archive; exclusion changes never rewrite source observations. + if current == 14 || current == 15 { structural, e := inspectStructuralCompatibilityMigrations(ctx, s, current, inspectPRDetailSchema(ctx, s)) if e != nil { return e @@ -296,7 +296,7 @@ func (s *Store) migrate(ctx context.Context) error { if e != nil { return e } - if len(structural) == 1 && structural[0] == "schema_version_14_to_15" && converged { + if len(structural) == 1 && structural[0] == fmt.Sprintf("schema_version_%d_to_%d", current, schemaVersion) && converged { // Some v14 archives predate the optional analytics extension. These // additive tables/columns are cheap to ensure and require no row scan. if e = s.ensureAnalyticsSourceSchema(ctx); e != nil { @@ -305,6 +305,9 @@ func (s *Store) migrate(ctx context.Context) error { if e = s.ensureAnalyticsIntegritySchema(ctx); e != nil { return e } + if e = s.ensureThreadExclusionsSchema(ctx); e != nil { + return e + } _, e = s.db.ExecContext(ctx, fmt.Sprintf("PRAGMA user_version=%d", schemaVersion)) return e } @@ -349,6 +352,9 @@ func (s *Store) migrate(ctx context.Context) error { if err := s.ensureAnalyticsIntegritySchema(ctx); err != nil { return err } + if err := s.ensureThreadExclusionsSchema(ctx); err != nil { + return err + } if err := s.ensureCanonicalObservationTables(ctx); err != nil { return err } diff --git a/internal/store/thread_exclusions.go b/internal/store/thread_exclusions.go new file mode 100644 index 00000000..f440ee1a --- /dev/null +++ b/internal/store/thread_exclusions.go @@ -0,0 +1,504 @@ +package store + +import ( + "context" + "crypto/sha256" + "database/sql" + "encoding/json" + "errors" + "fmt" + "path/filepath" + "sort" + "strings" + "time" +) + +var ErrThreadExcluded = errors.New("thread is excluded by owner policy") + +func (s *Store) ensureThreadExclusionsSchema(ctx context.Context) error { + _, err := s.q().ExecContext(ctx, ` +CREATE TABLE IF NOT EXISTS thread_exclusions( + repository TEXT NOT NULL,number INTEGER NOT NULL CHECK(number>0),kind TEXT NOT NULL, + original_thread_id INTEGER NOT NULL,github_id TEXT NOT NULL,excluded_at TEXT NOT NULL, + reason TEXT NOT NULL CHECK(reason='owner_requested'),plan_id TEXT NOT NULL, + PRIMARY KEY(repository,number)); +CREATE TABLE IF NOT EXISTS thread_excluded_nodes( + node_id TEXT PRIMARY KEY,repository TEXT NOT NULL,number INTEGER NOT NULL, + FOREIGN KEY(repository,number) REFERENCES thread_exclusions(repository,number)); +`) + if err != nil { + return err + } + + if s.hasTable(ctx, "analytics_review_state_coverage") && !s.hasColumn(ctx, "analytics_review_state_coverage", "owner_excluded_items") { + _, err = s.q().ExecContext(ctx, "ALTER TABLE analytics_review_state_coverage ADD COLUMN owner_excluded_items INTEGER NOT NULL DEFAULT 0") + } + return err +} + +func (s *Store) ensureThreadExclusionGuards(ctx context.Context) error { + _, err := s.q().ExecContext(ctx, ` +CREATE TRIGGER IF NOT EXISTS owner_excluded_thread_insert BEFORE INSERT ON threads + WHEN EXISTS(SELECT 1 FROM thread_exclusions e JOIN repositories r ON lower(r.full_name)=e.repository WHERE r.id=NEW.repo_id AND e.number=NEW.number) + BEGIN SELECT RAISE(ABORT,'owner_excluded_thread'); END; +CREATE TRIGGER IF NOT EXISTS owner_excluded_thread_update BEFORE UPDATE ON threads + WHEN EXISTS(SELECT 1 FROM thread_exclusions e JOIN repositories r ON lower(r.full_name)=e.repository WHERE r.id=NEW.repo_id AND e.number=NEW.number) + BEGIN SELECT RAISE(ABORT,'owner_excluded_thread'); END; +`) + if err != nil { + return err + } + for _, guard := range []struct{ table, name, condition, message string }{ + {"analytics_fetch_attempts", "attempt", "EXISTS(SELECT 1 FROM thread_exclusions WHERE repository=lower(NEW.repository) AND number=NEW.number)", "owner_excluded_thread"}, + {"sync_attempt_failures", "failure", "EXISTS(SELECT 1 FROM thread_exclusions e JOIN repositories r ON lower(r.full_name)=e.repository WHERE r.id=NEW.repo_id AND e.number=NEW.number)", "owner_excluded_thread"}, + {"analytics_retries", "retry", "EXISTS(SELECT 1 FROM thread_exclusions WHERE repository=lower(NEW.repository) AND number=NEW.number)", "owner_excluded_thread"}, + {"actor_identity_evidence", "identity", "EXISTS(SELECT 1 FROM thread_excluded_nodes WHERE node_id=NEW.node_id)", "owner_excluded_node"}, + {"analytics_pending_nodes", "pending_node", "EXISTS(SELECT 1 FROM thread_excluded_nodes WHERE node_id=NEW.node_id)", "owner_excluded_node"}, + } { + if !s.hasTable(ctx, guard.table) { + continue + } + for _, event := range []string{"INSERT", "UPDATE"} { + if _, err = s.q().ExecContext(ctx, "CREATE TRIGGER IF NOT EXISTS owner_excluded_"+guard.name+"_"+event+" BEFORE "+event+" ON "+guard.table+" WHEN "+guard.condition+" BEGIN SELECT RAISE(ABORT,'"+guard.message+"'); END"); err != nil { + return err + } + } + } + return nil +} + +func (s *Store) ThreadExcluded(ctx context.Context, repository string, number int) (bool, error) { + var excluded bool + err := s.q().QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM thread_exclusions WHERE repository=lower(?) AND number=?)", repository, number).Scan(&excluded) + return excluded, err +} +func (s *Store) NodeExcluded(ctx context.Context, id string) (bool, error) { + var excluded bool + err := s.q().QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM thread_excluded_nodes WHERE node_id=?)", id).Scan(&excluded) + return excluded, err +} +func (s *Store) FilterExcludedNumbers(ctx context.Context, repository string, numbers []int) ([]int, error) { + kept := make([]int, 0, len(numbers)) + for _, n := range numbers { + excluded, err := s.ThreadExcluded(ctx, repository, n) + if err != nil { + return nil, err + } + if !excluded { + kept = append(kept, n) + } + } + return kept, nil +} +func (s *Store) OwnerExcludedCount(ctx context.Context, repository string, reviewOnly bool) (int, error) { + var n int + err := s.q().QueryRowContext(ctx, "SELECT count(*) FROM thread_exclusions WHERE repository=lower(?) AND (?=0 OR kind='pull_request')", repository, boolInt(reviewOnly)).Scan(&n) + return n, err +} + +type PurgeTarget struct { + Number int `json:"number"` + ThreadID int64 `json:"thread_id"` + Kind string `json:"kind"` + GitHubID string `json:"github_id"` + Nodes []string `json:"content_node_ids"` +} +type ThreadPurgePlan struct { + Repository string `json:"repository"` + RepoID int64 `json:"repo_id"` + Targets []PurgeTarget `json:"targets"` + Counts map[string]int `json:"counts"` + PlanID string `json:"plan_id"` + Archive string `json:"archive"` + Applied bool `json:"applied"` + Reason string `json:"reason"` +} + +// Only source-owned relationships are traversed. Shared cluster descriptions +// need their own owner repair and are refused instead of discarding peer data. +func purgeRelations(ids string) map[string]string { + rev := "select id from thread_revisions where thread_id in (" + ids + ")" + snapshots := "select id from thread_code_snapshots where thread_revision_id in (" + rev + ")" + out := map[string]string{} + for _, table := range []string{"threads", "comments", "documents", "document_embeddings", "document_summaries", "pull_request_details", "pull_request_files", "pull_request_commits", "pull_request_checks", "pull_request_review_threads", "pull_request_review_thread_revisions", "pull_request_review_thread_syncs", "thread_revisions", "thread_vectors", "thread_child_observation_memberships", "thread_child_observation_reservations", "cluster_members", "cluster_memberships", "cluster_overrides"} { + col := "thread_id" + if table == "threads" { + col = "id" + } + out[table] = col + " in (" + ids + ")" + } + out["comment_revisions"] = "comment_id in (select id from comments where thread_id in (" + ids + "))" + for _, table := range []string{"thread_code_snapshots", "thread_fingerprints", "thread_key_summaries"} { + out[table] = "thread_revision_id in (" + rev + ")" + } + for _, table := range []string{"thread_changed_files", "thread_hunk_signatures"} { + out[table] = "snapshot_id in (" + snapshots + ")" + } + out["similarity_edges"] = "left_thread_id in (" + ids + ") or right_thread_id in (" + ids + ")" + return out +} + +var purgeBlobColumns = map[string][]string{"comments": {"raw_json_blob_id"}, "thread_revisions": {"raw_json_blob_id"}, "thread_changed_files": {"patch_blob_id"}, "thread_code_snapshots": {"raw_diff_blob_id"}} + +func (s *Store) PlanThreadPurge(ctx context.Context, repository string, numbers []int) (plan ThreadPurgePlan, err error) { + err = s.WithTx(ctx, func(tx *Store) error { + plan, err = tx.planThreadPurge(ctx, repository, numbers) + return err + }) + return plan, err +} + +func (s *Store) planThreadPurge(ctx context.Context, repository string, numbers []int) (ThreadPurgePlan, error) { + p := ThreadPurgePlan{Repository: strings.ToLower(strings.TrimSpace(repository)), Reason: "owner_requested", Counts: map[string]int{}, Targets: []PurgeTarget{}} + if s.hasTable(ctx, "portable_metadata") { + return p, fmt.Errorf("purge requires a native local archive") + } + archive, err := filepath.Abs(s.path) + if err != nil { + return p, err + } + p.Archive, err = filepath.EvalSymlinks(archive) + if err != nil { + return p, err + } + + if len(numbers) == 0 || len(numbers) > 100 { + return p, fmt.Errorf("purge requires 1..100 explicit numbers") + } + var canonical string + var matches int + if err := s.q().QueryRowContext(ctx, "SELECT count(*),coalesce(min(full_name),'') FROM repositories WHERE lower(full_name)=?", p.Repository).Scan(&matches, &canonical); err != nil { + return p, err + } + if matches != 1 { + return p, fmt.Errorf("purge repository must match exactly one local identity") + } + repo, err := s.RepositoryByFullName(ctx, canonical) + if err != nil { + return p, err + } + p.RepoID = repo.ID + seen := map[int]bool{} + var ids []string + for _, n := range numbers { + if n <= 0 || seen[n] { + return p, fmt.Errorf("purge numbers must be positive and unique") + } + seen[n] = true + t := PurgeTarget{Number: n, Nodes: []string{}} + var matches int + if err := s.q().QueryRowContext(ctx, "SELECT count(*) FROM threads WHERE repo_id=? AND number=?", repo.ID, n).Scan(&matches); err != nil { + return p, err + } + if matches > 1 { + return p, fmt.Errorf("ambiguous native thread number #%d", n) + } + err := s.q().QueryRowContext(ctx, "SELECT id,kind,github_id FROM threads WHERE repo_id=? AND number=?", repo.ID, n).Scan(&t.ThreadID, &t.Kind, &t.GitHubID) + if err != nil { + return p, fmt.Errorf("purge target #%d: %w", n, err) + } + ids = append(ids, fmt.Sprint(t.ThreadID)) + // historyProjection normalizes GraphQL id to node_id; top-level id + // is the REST/fullDatabaseId and must not become a node exclusion. + nodeQueries := []string{} + for _, rel := range []struct{ table, predicate string }{ + {"threads", fmt.Sprintf("id=%d", t.ThreadID)}, + {"comments", fmt.Sprintf("thread_id=%d", t.ThreadID)}, + {"comment_revisions", fmt.Sprintf("comment_id IN(SELECT id FROM comments WHERE thread_id=%d)", t.ThreadID)}, + {"thread_revisions", fmt.Sprintf("thread_id=%d", t.ThreadID)}, + } { + if s.hasColumn(ctx, rel.table, "raw_json") { + nodeQueries = append(nodeQueries, "SELECT json_extract(CASE WHEN json_valid(raw_json) THEN raw_json ELSE '{}' END,'$.node_id') node_id FROM "+rel.table+" WHERE "+rel.predicate) + } + } + if len(nodeQueries) > 0 { + rows, e := s.q().QueryContext(ctx, "SELECT DISTINCT node_id FROM ("+strings.Join(nodeQueries, " UNION ALL ")+") WHERE node_id IS NOT NULL AND node_id<>'' ORDER BY node_id") + if e != nil { + return p, e + } + for rows.Next() { + var node string + if e = rows.Scan(&node); e != nil { + rows.Close() + return p, e + } + t.Nodes = append(t.Nodes, node) + } + e = rows.Err() + rows.Close() + if e != nil { + return p, e + } + } + p.Targets = append(p.Targets, t) + } + sort.Slice(p.Targets, func(i, j int) bool { return p.Targets[i].Number < p.Targets[j].Number }) + idSQL := strings.Join(ids, ",") + relations := purgeRelations(idSQL) + for _, table := range []string{"clusters", "cluster_groups"} { + if !s.hasTable(ctx, table) { + continue + } + var n int + if err = s.q().QueryRowContext(ctx, "SELECT count(*) FROM "+table+" WHERE representative_thread_id in ("+idSQL+")").Scan(&n); err != nil { + return p, err + } + if n > 0 { + return p, fmt.Errorf("purge requires a separate cluster representative repair") + } + } + for _, table := range []string{"cluster_members", "cluster_memberships", "cluster_overrides"} { + var count int + if err = s.q().QueryRowContext(ctx, "SELECT count(*) FROM "+table+" WHERE "+relations[table]).Scan(&count); err != nil { + return p, err + } + if count > 0 { + return p, fmt.Errorf("purge requires a separate shared cluster repair") + } + } + // This deliberately bounded command does not resolve ownership inside + // retained workflow payloads. Refuse the repository surface rather than + // guessing from today's PR head or deleting a shared run. + if s.hasTable(ctx, "github_workflow_runs") { + var exists bool + if err = s.q().QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM github_workflow_runs WHERE repo_id=?)", p.RepoID).Scan(&exists); err != nil { + return p, err + } + if exists { + return p, fmt.Errorf("purge requires a separate repository workflow-snapshot repair") + } + } + heads := []string{} + if s.hasTable(ctx, "pull_request_details") { + heads = append(heads, "SELECT head_sha FROM pull_request_details WHERE thread_id IN("+idSQL+")") + } + if s.hasTable(ctx, "thread_code_snapshots") { + heads = append(heads, "SELECT head_sha FROM thread_code_snapshots WHERE "+relations["thread_code_snapshots"]) + } + for _, table := range []string{"threads", "thread_revisions"} { + if !s.hasColumn(ctx, table, "raw_json") { + continue + } + for _, path := range []string{"$._graphql.headRefOid", "$.head.sha"} { + heads = append(heads, "SELECT json_extract(CASE WHEN json_valid(raw_json) THEN raw_json ELSE '{}' END,'"+path+"') FROM "+table+" WHERE "+relations[table]) + } + } + if s.hasTable(ctx, "workflow_run_observation_reservations") && len(heads) > 0 { + var n int + if err = s.q().QueryRowContext(ctx, "SELECT count(*) FROM workflow_run_observation_reservations WHERE repo_id=? AND head_sha IN("+strings.Join(heads, " UNION ")+")", p.RepoID).Scan(&n); err != nil { + return p, err + } + if n > 0 { + return p, fmt.Errorf("purge requires a separate linked workflow-reservation repair") + } + p.Counts["workflow_run_observation_reservations"] = 0 + } + // Blob-backed payloads need an identity/ownership-aware repair of their + // own. Do not silently leave their nodes behind or read unknown blob prose. + for table, cols := range purgeBlobColumns { + for _, col := range cols { + if !s.hasColumn(ctx, table, col) { + continue + } + var exists bool + if err = s.q().QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM "+table+" WHERE ("+relations[table]+") AND "+col+" IS NOT NULL)").Scan(&exists); err != nil { + return p, err + } + if exists { + return p, fmt.Errorf("purge requires a separate blob-backed payload repair") + } + } + } + + var numeric []string + var nodes []string + for _, t := range p.Targets { + numeric = append(numeric, fmt.Sprint(t.Number)) + for _, node := range t.Nodes { + nodes = append(nodes, "'"+strings.ReplaceAll(node, "'", "''")+"'") + } + } + + // Node identity belongs to one conversation. Refuse corrupt/shared native + // keys instead of deleting a peer's evidence or suppressing its enrichment. + if len(nodes) > 0 { + for _, table := range []string{"threads", "comments", "comment_revisions"} { + var shared bool + if err := s.q().QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM "+table+" WHERE NOT ("+relations[table]+") AND json_extract(CASE WHEN json_valid(raw_json) THEN raw_json ELSE '{}' END,'$.node_id') IN("+strings.Join(nodes, ",")+"))").Scan(&shared); err != nil { + return p, err + } + if shared { + return p, fmt.Errorf("purge refuses shared content node identities") + } + } + } + for _, table := range []string{"analytics_retries", "analytics_fetch_attempts"} { + relations[table] = fmt.Sprintf("lower(repository)='%s' AND number IN(%s)", strings.ReplaceAll(p.Repository, "'", "''"), strings.Join(numeric, ",")) + } + relations["sync_attempt_failures"] = fmt.Sprintf("repo_id=%d AND number IN(%s)", p.RepoID, strings.Join(numeric, ",")) + for _, table := range []string{"actor_identity_evidence", "analytics_pending_nodes"} { + relations[table] = "node_id IN(" + strings.Join(nodes, ",") + ")" + } + + // A durable downstream exclusion is keyed by native row ID. Refuse a tail + // deletion that SQLite could later reuse; never create dummy rows or IDs. + for table, predicate := range relations { + var singleID bool + if err := s.q().QueryRowContext(ctx, "SELECT count(*)=1 AND min(name)='id' AND min(upper(type))='INTEGER' FROM pragma_table_info(?) WHERE pk>0", table).Scan(&singleID); err != nil { + return p, err + } + if singleID { + if err := s.guardPurgeIDTail(ctx, table, predicate); err != nil { + return p, err + } + } + } + // Hash complete selected rows, not just counts: edits between preview and + // apply must require a new plan. Sorted tables and rowids make it stable. + tables := make([]string, 0, len(relations)) + for table := range relations { + tables = append(tables, table) + } + sort.Strings(tables) + digest := sha256.New() + enc := json.NewEncoder(digest) + for _, table := range tables { + p.Counts[table] = 0 + if !s.hasTable(ctx, table) { + continue + } + rows, err := s.q().QueryContext(ctx, "SELECT * FROM "+table+" WHERE "+relations[table]+" ORDER BY rowid") + if err != nil { + return p, err + } + columns, err := rows.Columns() + if err != nil { + rows.Close() + return p, err + } + values := make([]any, len(columns)) + pointers := make([]any, len(columns)) + for i := range values { + pointers[i] = &values[i] + } + p.Counts[table] = 0 + for rows.Next() { + if err = rows.Scan(pointers...); err != nil { + break + } + if err = enc.Encode([]any{table, columns, values}); err != nil { + break + } + p.Counts[table]++ + } + rowErr := rows.Err() + rows.Close() + if err != nil { + return p, err + } + if rowErr != nil { + return p, rowErr + } + } + if err = enc.Encode(p); err != nil { + return p, err + } + p.PlanID = fmt.Sprintf("%x", digest.Sum(nil)) + + return p, nil +} + +func (s *Store) guardPurgeIDTail(ctx context.Context, table, predicate string) error { + var selected, maximum sql.NullInt64 + if err := s.q().QueryRowContext(ctx, "SELECT max(id) FROM "+table+" WHERE "+predicate).Scan(&selected); err != nil { + return err + } + if !selected.Valid { + return nil + } + if err := s.q().QueryRowContext(ctx, "SELECT max(id) FROM "+table).Scan(&maximum); err != nil { + return err + } + if selected.Int64 == maximum.Int64 { + return fmt.Errorf("purge would permit native ID reuse in %s; a retained higher ID is required", table) + } + return nil +} + +// PurgeThreads is owner-directed removal, not a provider deletion or successful +// fetch. Caller holds runner.lock; a single transaction fences all native writes. +func (s *Store) PurgeThreads(ctx context.Context, repository string, numbers []int, planID string) (ThreadPurgePlan, error) { + var result ThreadPurgePlan + if len(planID) != 64 { + return result, fmt.Errorf("purge requires the plan_id from a fresh preview") + } + err := s.WithTx(ctx, func(tx *Store) error { + if err := tx.ensureThreadExclusionsSchema(ctx); err != nil { + return err + } + if err := tx.ensureThreadExclusionGuards(ctx); err != nil { + return err + } + p, err := tx.planThreadPurge(ctx, repository, numbers) + if err != nil { + return err + } + if p.PlanID != planID { + return fmt.Errorf("purge plan changed; preview again before applying") + } + var enabled int + if err = tx.q().QueryRowContext(ctx, "PRAGMA foreign_keys").Scan(&enabled); err != nil { + return err + } + if enabled != 1 { + return fmt.Errorf("purge requires foreign keys") + } + at := time.Now().UTC().Format(time.RFC3339Nano) + for _, t := range p.Targets { + if _, err = tx.q().ExecContext(ctx, `INSERT INTO thread_exclusions(repository,number,kind,original_thread_id,github_id,excluded_at,reason,plan_id) VALUES(?,?,?,?,?,?,'owner_requested',?)`, p.Repository, t.Number, t.Kind, t.ThreadID, t.GitHubID, at, planID); err != nil { + return err + } + for _, node := range t.Nodes { + if _, err = tx.q().ExecContext(ctx, "INSERT INTO thread_excluded_nodes(node_id,repository,number) VALUES(?,?,?)", node, p.Repository, t.Number); err != nil { + return err + } + for _, table := range []string{"actor_identity_evidence", "analytics_pending_nodes"} { + if !tx.hasTable(ctx, table) { + continue + } + if _, err = tx.q().ExecContext(ctx, "DELETE FROM "+table+" WHERE node_id=?", node); err != nil { + return err + } + } + } + for _, table := range []string{"analytics_retries", "analytics_fetch_attempts"} { + if !tx.hasTable(ctx, table) { + continue + } + if _, err = tx.q().ExecContext(ctx, "DELETE FROM "+table+" WHERE lower(repository)=? AND number=?", p.Repository, t.Number); err != nil { + return err + } + } + if tx.hasTable(ctx, "sync_attempt_failures") { + if _, err = tx.q().ExecContext(ctx, "DELETE FROM sync_attempt_failures WHERE repo_id=? AND number=?", p.RepoID, t.Number); err != nil { + return err + } + } + if _, err = tx.q().ExecContext(ctx, "DELETE FROM threads WHERE id=?", t.ThreadID); err != nil { + return err + } + } + // Never turn removal of unavailable work into claimed provider completeness. + if tx.hasTable(ctx, "analytics_review_state_coverage") { + if _, err = tx.q().ExecContext(ctx, `UPDATE analytics_review_state_coverage SET pending_items=(SELECT count(*) FROM analytics_retries WHERE lower(repository)=? AND operation='review_state' AND resolved_at IS NULL),owner_excluded_items=(SELECT count(*) FROM thread_exclusions WHERE repository=? AND kind='pull_request'),complete=0 WHERE lower(repository)=?`, p.Repository, p.Repository, p.Repository); err != nil { + return err + } + } + result = p + return nil + }) + if err == nil { + result.Applied = true + } + return result, err +} diff --git a/internal/store/thread_exclusions_test.go b/internal/store/thread_exclusions_test.go new file mode 100644 index 00000000..91d3e5e0 --- /dev/null +++ b/internal/store/thread_exclusions_test.go @@ -0,0 +1,484 @@ +package store + +import ( + "context" + "errors" + "fmt" + "path/filepath" + "strings" + "testing" +) + +func exclusionFixture(t *testing.T) *Store { + t.Helper() + s, err := Open(context.Background(), filepath.Join(t.TempDir(), "archive.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { s.Close() }) + _, err = s.DB().Exec(` +INSERT INTO repositories(id,owner,name,full_name,github_repo_id,raw_json,updated_at) VALUES(1,'fixture','repo','fixture/repo','100','{}','2026-01-01'); +INSERT INTO threads(id,repo_id,github_id,number,kind,state,title,body,html_url,labels_json,assignees_json,raw_json,content_hash,updated_at) VALUES + (1,1,'gh10',10,'pull_request','closed','remove','OWNER_REMOVAL_SENTINEL','https://github.com/fixture/repo/pull/10','[]','[]','{"id":1010,"node_id":"node10","_graphql":{"id":"node10"}}','hash10','2026-01-01'), + (2,1,'gh20',20,'pull_request','open','keep','UNCHANGED_BODY','https://github.com/fixture/repo/pull/20','[]','[]','{"node_id":"node20"}','hash20','2026-01-01'); +INSERT INTO comments(id,thread_id,github_id,comment_type,body,raw_json) VALUES + (1,1,'comment10','issue_comment','removed comment','{"node_id":"comment-node10"}'), + (2,2,'comment20','issue_comment','kept comment','{"node_id":"comment-node20"}'); +INSERT INTO comment_revisions(id,comment_id,body,raw_json,recorded_at) VALUES + (1,1,'removed old comment','{"node_id":"comment-node10"}','2026-01-01'), + (2,2,'kept old comment','{"node_id":"comment-node20"}','2026-01-01'); +INSERT INTO documents(id,thread_id,title,body,raw_text,dedupe_text,updated_at) VALUES + (1,1,'remove','OWNER_REMOVAL_SENTINEL','OWNER_REMOVAL_SENTINEL','remove','2026-01-01'), + (2,2,'keep','UNCHANGED_BODY','UNCHANGED_BODY','keep','2026-01-01'); +INSERT INTO pull_request_details(thread_id,repo_id,number,raw_json,fetched_at,updated_at) VALUES(1,1,10,'{"fixture":"removed"}','2026-01-01','2026-01-01'),(2,1,20,'{"fixture":"kept"}','2026-01-01','2026-01-01'); +INSERT INTO thread_revisions(id,thread_id,content_hash,title_hash,body_hash,labels_hash,created_at) VALUES + (1,1,'remove','t','b','l','2026-01-01'),(2,2,'keep','t','b','l','2026-01-01'); +INSERT INTO thread_key_summaries(id,thread_revision_id,summary_kind,prompt_version,provider,model,input_hash,output_hash,key_text,created_at) VALUES + (1,1,'key','v1','fixture','fixture','i','o','removed summary','2026-01-01'),(2,2,'key','v1','fixture','fixture','i','o','kept summary','2026-01-01'); +INSERT INTO sync_attempt_failures(id,repo_id,thread_id,number,operation,error_class,error_message,first_seen_at,last_seen_at) VALUES + (1,1,1,10,'issue','fixture','removed error','2026-01-01','2026-01-01'),(2,1,2,20,'issue','fixture','kept error','2026-01-01','2026-01-01'); +INSERT INTO actor_identity_evidence VALUES('node10','shared-actor','fixture','User','2026-01-01','{}'),('comment-node10','shared-actor','fixture','User','2026-01-01','{}'); +INSERT INTO actor_profiles(node_id,login,actor_type,observed_at,raw_json) VALUES('shared-actor','fixture','User','2026-01-01','{}'); +INSERT INTO analytics_pending_nodes VALUES('node10','identity'),('comment-node10','identity'),('shared-actor','profile'); +`) + if err != nil { + t.Fatal(err) + } + ctx := context.Background() + at := "2026-01-01T00:00:00Z" + if err = s.SaveAnalyticsCoverage(ctx, "fixture/repo", at, 0, 2); err != nil { + t.Fatal(err) + } + for _, n := range []int{10, 20} { + if err = s.RecordAnalyticsAttempt(ctx, AnalyticsAttempt{Repository: "fixture/repo", Number: n, Operation: "review_state", StartedAt: at, FinishedAt: at, Status: "failed", ErrorClass: "partial_response", Evidence: []byte(`{"numbers":[10,20],"body_evidence":{"sha256":"fixture-hash","bytes":10}}`)}); err != nil { + t.Fatal(err) + } + } + if err = s.SaveReviewStateCoverage(ctx, "fixture/repo", ReviewStateRecovery{Done: true, Cursor: 2, Ceiling: 2, Scanned: 2, Queued: 2}); err != nil { + t.Fatal(err) + } + return s +} + +func TestOwnerPurgeRemovesExactClosureAndPreventsReintroduction(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + plan, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil { + t.Fatal(err) + } + if plan.Applied || len(plan.Targets) != 1 || len(plan.Targets[0].Nodes) != 2 || plan.Counts["comments"] != 1 || plan.Counts["comment_revisions"] != 1 { + t.Fatalf("plan=%+v", plan) + } + var retained string + s.DB().QueryRow("SELECT raw_json FROM pull_request_details WHERE thread_id=2").Scan(&retained) + result, err := applyPurge(t, s, ctx, "fixture/repo", []int{10}, "owner-request-fixture") + if err != nil || !result.Applied { + t.Fatalf("%+v %v", result, err) + } + for _, table := range []string{"threads", "comments", "comment_revisions", "documents", "pull_request_details", "thread_revisions", "thread_key_summaries", "sync_attempt_failures"} { + var n int + if err = s.DB().QueryRow("SELECT count(*) FROM " + table).Scan(&n); err != nil || n != 1 { + t.Fatalf("%s: %d %v", table, n, err) + } + } + var after string + s.DB().QueryRow("SELECT raw_json FROM pull_request_details WHERE thread_id=2").Scan(&after) + if after != retained { + t.Fatal("peer evidence changed") + } + var fts int + if err = s.DB().QueryRow("SELECT count(*) FROM documents_fts WHERE documents_fts MATCH 'OWNER_REMOVAL_SENTINEL'").Scan(&fts); err != nil || fts != 0 { + t.Fatalf("FTS leaked removed document: %d %v", fts, err) + } + var profiles int + s.DB().QueryRow("SELECT count(*) FROM actor_profiles").Scan(&profiles) + if profiles != 1 { + t.Fatal("shared actor removed") + } + var remaining string + s.DB().QueryRow("SELECT evidence_json FROM analytics_fetch_attempts WHERE number=20").Scan(&remaining) + if !strings.Contains(remaining, `"numbers":[10,20]`) { + t.Fatal("unrelated mixed receipt changed") + } + if excluded, err := s.ThreadExcluded(ctx, "FIXTURE/REPO", 10); err != nil || !excluded { + t.Fatal("missing permanent policy") + } + if excluded, err := s.NodeExcluded(ctx, "1010"); err != nil || excluded { + t.Fatal("REST database ID mistaken for GraphQL node", err) + } + if _, err = s.UpsertThread(ctx, Thread{RepoID: 1, Number: 10, Kind: "pull_request", GitHubID: "gh10"}); !errors.Is(err, ErrThreadExcluded) { + t.Fatalf("reinsertion allowed: %v", err) + } + if err = s.SaveActorEvidence(ctx, []map[string]any{{"id": "node10", "author": map[string]any{"id": "shared-actor"}}}, "2026-01-02T00:00:00Z"); err != nil { + t.Fatal(err) + } + if err = s.RecordAnalyticsAttempt(ctx, AnalyticsAttempt{Repository: "fixture/repo", Number: 10, Operation: "review_state", StartedAt: "2026-01-02T00:00:00Z", FinishedAt: "2026-01-02T00:00:00Z", Status: "success", Evidence: []byte(`{}`)}); err != nil { + t.Fatal(err) + } + for _, q := range []string{"SELECT count(*) FROM analytics_fetch_attempts WHERE number=10", "SELECT count(*) FROM analytics_retries WHERE number=10", "SELECT count(*) FROM actor_identity_evidence WHERE node_id='node10'", "SELECT count(*) FROM analytics_pending_nodes WHERE node_id='node10'"} { + var n int + s.DB().QueryRow(q).Scan(&n) + if n != 0 { + t.Fatal("excluded evidence recreated", q, n) + } + } + if _, err = s.DB().Exec("INSERT INTO analytics_pending_nodes VALUES('node10','identity')"); err == nil { + t.Fatal("SQL guard bypass") + } + // The remaining ordinary retry can resolve, but owner removal is not provider completeness. + s.DB().Exec("UPDATE analytics_retries SET resolved_at='2026-01-02T00:00:00Z' WHERE number=20") + if err = s.SaveReviewStateCoverage(ctx, "fixture/repo", ReviewStateRecovery{Done: true}); err != nil { + t.Fatal(err) + } + var complete, pending, excluded, core int + s.DB().QueryRow("SELECT complete,pending_items,owner_excluded_items FROM analytics_review_state_coverage").Scan(&complete, &pending, &excluded) + s.DB().QueryRow("SELECT complete FROM analytics_coverage").Scan(&core) + if complete != 0 || pending != 0 || excluded != 1 || core != 1 { + t.Fatalf("false recovery: review=%d pending=%d excluded=%d core=%d", complete, pending, excluded, core) + } + if _, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}); err == nil { + t.Fatal("removed target should require a new selection") + } + var original string + if err := s.DB().QueryRow("SELECT plan_id FROM thread_exclusions").Scan(&original); err != nil || original != plan.PlanID { + t.Fatal("plan identity not retained", err) + } + +} + +func TestOwnerPurgeRejectsTailAndRollsBack(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + if _, err := applyPurge(t, s, ctx, "fixture/repo", []int{20}, "tail"); err == nil || !strings.Contains(err.Error(), "ID reuse") { + t.Fatalf("tail admitted: %v", err) + } + var n int + s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&n) + if n != 0 { + t.Fatal("tail failure wrote policy") + } + id, err := s.UpsertThread(ctx, Thread{RepoID: 1, Number: 30, Kind: "issue", GitHubID: "gh30", Title: "new", RawJSON: "{}", LabelsJSON: "[]", AssigneesJSON: "[]", UpdatedAt: "2026-01-02T00:00:00Z"}) + if err != nil || id <= 2 { + t.Fatalf("native ID reused: %d %v", id, err) + } + if _, err = s.PlanThreadPurge(ctx, "fixture/repo", []int{20}); err == nil || !strings.Contains(err.Error(), "ID reuse") { + t.Fatal("child tail admitted", err) + } + _, err = s.DB().Exec(`CREATE TRIGGER deny_fixture_delete BEFORE DELETE ON comments WHEN OLD.id=1 BEGIN SELECT RAISE(ABORT,'fixture interruption'); END`) + if err != nil { + t.Fatal(err) + } + if _, err = applyPurge(t, s, ctx, "fixture/repo", []int{10}, "rollback"); err == nil { + t.Fatal("failure accepted") + } + for _, q := range []string{"SELECT count(*) FROM threads WHERE id=1", "SELECT count(*) FROM actor_identity_evidence WHERE node_id='node10'", "SELECT count(*) FROM analytics_retries WHERE number=10"} { + s.DB().QueryRow(q).Scan(&n) + if n != 1 { + t.Fatal("partial purge", q, n) + } + } + s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&n) + if n != 0 { + t.Fatal("rollback left exclusion") + } + for _, numbers := range [][]int{nil, {0}, {10, 10}, {999}} { + if _, err = s.PlanThreadPurge(ctx, "fixture/repo", numbers); err == nil { + t.Fatal("invalid selection accepted", numbers) + } + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10}, ""); err == nil { + t.Fatal("missing owner request accepted") + } +} + +func TestOwnerExclusionV15MigrationIsAdditive(t *testing.T) { + s := exclusionFixture(t) + path := s.Path() + _, err := s.DB().Exec(`DROP TABLE thread_excluded_nodes;DROP TABLE thread_exclusions;ALTER TABLE analytics_review_state_coverage DROP COLUMN owner_excluded_items;PRAGMA user_version=15;`) + if err != nil { + t.Fatal(err) + } + if err = s.markObservationSchemaConverged(context.Background()); err != nil { + t.Fatal(err) + } + s.Close() + s, err = Open(context.Background(), path) + if err != nil { + t.Fatal(err) + } + defer s.Close() + var version, count int + s.DB().QueryRow("pragma user_version").Scan(&version) + s.DB().QueryRow("select count(*) from comments").Scan(&count) + if version != 16 || count != 2 { + t.Fatalf("migration altered data: %d %d", version, count) + } +} + +func TestOwnerPurgeRefusesBlobBackedTargetsWithoutMutation(t *testing.T) { + for _, mode := range []string{"shared", "unshared", "external", "tail"} { + t.Run(mode, func(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + _, err := s.DB().Exec(`INSERT INTO blobs(id,sha256,media_type,size_bytes,storage_kind,inline_text,created_at) VALUES + (1,'one','application/json',20,'inline','removed-blob','2026-01-01'), + (2,'two','application/json',20,'inline','peer-blob','2026-01-01'); + UPDATE comments SET raw_json_blob_id=1 WHERE id=1;`) + if err != nil { + t.Fatal(err) + } + switch mode { + case "shared": + _, err = s.DB().Exec("UPDATE comments SET raw_json_blob_id=1 WHERE id=2") + case "external": + _, err = s.DB().Exec("UPDATE blobs SET storage_kind='file',storage_path='/fixture/private' WHERE id=1") + case "tail": + _, err = s.DB().Exec("DELETE FROM blobs WHERE id=2") + } + if err != nil { + t.Fatal(err) + } + result, err := applyPurge(t, s, ctx, "fixture/repo", []int{10}, "blob-fixture") + var count, threads int + s.DB().QueryRow("SELECT count(*) FROM blobs WHERE id=1").Scan(&count) + s.DB().QueryRow("SELECT count(*) FROM threads WHERE id=1").Scan(&threads) + if err == nil || result.Applied || !strings.Contains(err.Error(), "blob-backed") || count != 1 || threads != 1 { + t.Fatalf("unsupported blob target was altered: %+v %v count=%d threads=%d", result, err, count, threads) + } + var policy int + s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&policy) + if policy != 0 { + t.Fatal("refusal installed an exclusion") + } + + }) + } +} + +func TestOwnerPurgePublicationCannotDiscardPolicy(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + if _, err := applyPurge(t, s, ctx, "fixture/repo", []int{10}, "private-owner-reference"); err != nil { + t.Fatal(err) + } + if _, err := s.PrunePortablePayloads(ctx, PortablePruneOptions{BodyChars: 1}); err == nil || !strings.Contains(err.Error(), "local owner exclusions") { + t.Fatal("private policy publication admitted", err) + } + var body string + if err := s.DB().QueryRow("SELECT body FROM threads WHERE id=2").Scan(&body); err != nil || body != "UNCHANGED_BODY" { + t.Fatal("refusal modified retained data", err) + } +} + +func TestOwnerPurgeCanonicalRepositoryMatchesCaseVariantReceipts(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + _, err := s.DB().Exec("UPDATE repositories SET full_name='Fixture/Repo';UPDATE analytics_fetch_attempts SET repository='Fixture/Repo' WHERE number=10;UPDATE analytics_retries SET repository='Fixture/Repo' WHERE number=10;UPDATE analytics_review_state_coverage SET repository='Fixture/Repo'") + if err != nil { + t.Fatal(err) + } + result, err := applyPurge(t, s, ctx, "FIXTURE/REPO", []int{10}, "case-fixture") + if err != nil || result.Counts["analytics_fetch_attempts"] != 1 || result.Counts["analytics_retries"] != 1 { + t.Fatalf("variant not planned: %+v %v", result, err) + } + for _, table := range []string{"analytics_fetch_attempts", "analytics_retries"} { + var n int + if err = s.DB().QueryRow("SELECT count(*) FROM " + table + " WHERE number=10").Scan(&n); err != nil || n != 0 { + t.Fatal("variant retained", table, n, err) + } + } + var complete, excluded int + s.DB().QueryRow("SELECT complete,owner_excluded_items FROM analytics_review_state_coverage WHERE repository='Fixture/Repo'").Scan(&complete, &excluded) + if complete != 0 || excluded != 1 { + t.Fatal("variant coverage falsely complete", complete, excluded) + } + if err = s.RecordAnalyticsAttempt(ctx, AnalyticsAttempt{Repository: "Fixture/Repo", Number: 10, Operation: "review_state", StartedAt: "2026-01-01T00:00:00Z", FinishedAt: "2026-01-01T00:00:01Z", Status: "failed", ErrorClass: "partial_response", Evidence: []byte(`{}`)}); err != nil { + t.Fatal(err) + } + var n int + s.DB().QueryRow("SELECT count(*) FROM analytics_retries WHERE number=10").Scan(&n) + if n != 0 { + t.Fatal("variant requeued") + } +} + +func TestOwnerPurgeRefusesWorkflowEvidenceAndKeepsUnrelatedReservations(t *testing.T) { + for _, mode := range []string{"run", "current-head", "historical-head", "unrelated"} { + t.Run(mode, func(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + var err error + switch mode { + case "run": + _, err = s.DB().Exec("INSERT INTO github_workflow_runs(repo_id,run_id,head_sha,raw_json,fetched_at) VALUES(1,1,'older-sha','{}','2026-01-01')") + case "current-head": + _, err = s.DB().Exec("UPDATE pull_request_details SET head_sha='linked-sha' WHERE thread_id=1;INSERT INTO workflow_run_observation_reservations VALUES(1,'linked-sha','2026-01-01',1)") + case "historical-head": + _, err = s.DB().Exec(`UPDATE threads SET raw_json='{"node_id":"node10","_graphql":{"headRefOid":"older-sha"}}' WHERE id=1;INSERT INTO workflow_run_observation_reservations VALUES(1,'older-sha','2026-01-01',1)`) + case "unrelated": + _, err = s.DB().Exec("UPDATE pull_request_details SET head_sha='peer-sha' WHERE thread_id=2;INSERT INTO workflow_run_observation_reservations VALUES(1,'peer-sha','2026-01-01',1)") + } + if err != nil { + t.Fatal(err) + } + result, err := applyPurge(t, s, ctx, "fixture/repo", []int{10}, "workflow-fixture") + var n int + s.DB().QueryRow("SELECT count(*) FROM threads WHERE id=1").Scan(&n) + if mode == "unrelated" { + if err != nil || !result.Applied || n != 0 { + t.Fatalf("unrelated reservation blocked target: %+v %v", result, err) + } + s.DB().QueryRow("SELECT count(*) FROM workflow_run_observation_reservations WHERE head_sha='peer-sha'").Scan(&n) + if n != 1 { + t.Fatal("peer reservation lost") + } + } else { + if err == nil || !strings.Contains(err.Error(), "workflow") || result.Applied || n != 1 { + t.Fatalf("workflow evidence ignored: %+v %v count=%d", result, err, n) + } + s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&n) + if n != 0 { + t.Fatal("refusal wrote policy") + } + } + }) + } +} + +func applyPurge(t *testing.T, s *Store, ctx context.Context, repository string, numbers []int, _ string) (ThreadPurgePlan, error) { + t.Helper() + plan, err := s.PlanThreadPurge(ctx, repository, numbers) + if err != nil { + return plan, err + } + return s.PurgeThreads(ctx, repository, numbers, plan.PlanID) +} + +func TestOwnerPurgePlanIdentityRejectsStaleContent(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + plan, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil { + t.Fatal(err) + } + again, err := s.PlanThreadPurge(ctx, "FIXTURE/REPO", []int{10}) + if err != nil || plan.PlanID != again.PlanID || len(plan.PlanID) != 64 { + t.Fatal("unstable plan", err) + } + if _, err = s.DB().Exec("UPDATE comments SET body='changed after preview' WHERE id=1"); err != nil { + t.Fatal(err) + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10}, plan.PlanID); err == nil || !strings.Contains(err.Error(), "plan changed") { + t.Fatal("stale plan accepted", err) + } + var n int + if err = s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&n); err != nil || n != 0 { + t.Fatal("stale apply changed policy", err) + } + fresh, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil { + t.Fatal(err) + } + if _, err = s.DB().Exec("CREATE TRIGGER fail_mid_purge BEFORE DELETE ON comments BEGIN SELECT RAISE(ABORT,'injected mid-apply failure'); END"); err != nil { + t.Fatal(err) + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10}, fresh.PlanID); err == nil { + t.Fatal("failure injection ignored") + } + after, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil || after.PlanID != fresh.PlanID { + t.Fatal("rollback changed selected rows", err) + } + if err = s.DB().QueryRow("SELECT count(*) FROM thread_exclusions").Scan(&n); err != nil || n != 0 { + t.Fatal("rollback changed policy", err) + } +} + +func TestOwnerPurgeRefusesAmbiguousAndSharedTargets(t *testing.T) { + for _, mutation := range []string{ + "UPDATE comments SET raw_json='{\"node_id\":\"node10\"}' WHERE id=2", + "INSERT INTO repositories(owner,name,full_name,raw_json,updated_at) VALUES('Fixture','Repo','Fixture/Repo','{}','now')", + "INSERT INTO threads(repo_id,github_id,number,kind,state,title,html_url,labels_json,assignees_json,raw_json,content_hash,updated_at) VALUES(1,'other',10,'issue','open','','','[]','[]','{}','h','now')", + "INSERT INTO cluster_groups(repo_id,stable_key,stable_slug,status,representative_thread_id,created_at,updated_at) VALUES(1,'k','s','open',1,'now','now')", + "CREATE TABLE portable_metadata(fixture TEXT)", + } { + t.Run(mutation, func(t *testing.T) { + s := exclusionFixture(t) + if _, err := s.DB().Exec(mutation); err != nil { + t.Fatal(err) + } + if _, err := s.PlanThreadPurge(context.Background(), "fixture/repo", []int{10}); err == nil { + t.Fatal("unsafe selection accepted") + } + }) + } +} + +func TestOwnerPurgeMultipleTargetsAreOneTransaction(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + for _, n := range []int{30, 40} { + if _, err := s.UpsertThread(ctx, Thread{RepoID: 1, Number: n, GitHubID: fmt.Sprint(n), Kind: "issue", Title: "synthetic", RawJSON: "{}", LabelsJSON: "[]", AssigneesJSON: "[]", UpdatedAt: "2026-01-01T00:00:00Z"}); err != nil { + t.Fatal(err) + } + } + plan, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{30, 10}) + if err != nil { + t.Fatal(err) + } + ordered, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10, 30}) + if err != nil || ordered.PlanID != plan.PlanID || plan.Counts["threads"] != 2 || plan.Targets[0].Number != 10 || plan.Targets[1].Number != 30 { + t.Fatal("selection or order changed plan", err) + } + if _, err = s.DB().Exec("CREATE TRIGGER fail_second_target BEFORE DELETE ON threads WHEN OLD.number=30 BEGIN SELECT RAISE(ABORT,'second target failed'); END"); err != nil { + t.Fatal(err) + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10, 30}, plan.PlanID); err == nil { + t.Fatal("second-target failure ignored") + } + after, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10, 30}) + if err != nil || after.PlanID != plan.PlanID { + t.Fatal("transaction changed selected data", err) + } + if _, err = s.DB().Exec("DROP TRIGGER fail_second_target"); err != nil { + t.Fatal(err) + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10, 30}, plan.PlanID); err != nil { + t.Fatal(err) + } + var kept string + if err = s.DB().QueryRow("SELECT group_concat(number,',') FROM (SELECT number FROM threads ORDER BY number)").Scan(&kept); err != nil || kept != "20,40" { + t.Fatal("wrong retained set", kept, err) + } +} + +func TestOwnerPurgePreviewBeforeOptionalAnalyticsMigration(t *testing.T) { + s := exclusionFixture(t) + ctx := context.Background() + path := s.Path() + if _, err := s.DB().Exec(`DROP TABLE analytics_fetch_attempts;DROP TABLE analytics_retries;DROP TABLE actor_identity_evidence;DROP TABLE analytics_pending_nodes;PRAGMA user_version=14`); err != nil { + t.Fatal(err) + } + s.Close() + readonly, err := OpenReadOnly(ctx, path) + if err != nil { + t.Fatal(err) + } + plan, err := readonly.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil { + t.Fatal(err) + } + if readonly.hasTable(ctx, "analytics_retries") || plan.Counts["analytics_retries"] != 0 { + t.Fatal("preview migrated archive") + } + readonly.Close() + writable, err := Open(ctx, path) + if err != nil { + t.Fatal(err) + } + defer writable.Close() + if _, err := writable.PurgeThreads(ctx, "fixture/repo", []int{10}, plan.PlanID); err != nil { + t.Fatal("empty optional tables changed selection", err) + } +} diff --git a/internal/store/threads.go b/internal/store/threads.go index 7b1b2a32..45375ff5 100644 --- a/internal/store/threads.go +++ b/internal/store/threads.go @@ -81,6 +81,13 @@ func (s *Store) UpsertThreadObservation(ctx context.Context, thread Thread, opti } func (s *Store) upsertThreadObservation(ctx context.Context, thread Thread, options UpsertThreadOptions) (UpsertThreadResult, error) { + var excluded bool + if err := s.q().QueryRowContext(ctx, `SELECT EXISTS(SELECT 1 FROM thread_exclusions e JOIN repositories r ON lower(r.full_name)=e.repository WHERE r.id=? AND e.number=?)`, thread.RepoID, thread.Number).Scan(&excluded); err != nil { + return UpsertThreadResult{}, err + } + if excluded { + return UpsertThreadResult{}, ErrThreadExcluded + } if options.ObservationSequence <= 0 { sequence, err := s.NextThreadObservationSequence(ctx, thread.UpdatedAt) if err != nil { diff --git a/internal/syncer/graphql_history_test.go b/internal/syncer/graphql_history_test.go index e8faa959..4ac01adf 100644 --- a/internal/syncer/graphql_history_test.go +++ b/internal/syncer/graphql_history_test.go @@ -14,11 +14,15 @@ import ( type historyFixtureClient struct { *gh.Client - batch gh.HistoryBatch - err error + batch gh.HistoryBatch + err error + beforeFetch func() } func (f historyFixtureClient) FetchGraphQLHistory(context.Context, string, string, []int, gh.Reporter) (gh.HistoryBatch, error) { + if f.beforeFetch != nil { + f.beforeFetch() + } return f.batch, f.err } func TestGraphQLHistoryUsesNativeTransactionsAndPreservesLegacyIdentity(t *testing.T) { @@ -109,8 +113,9 @@ func TestGraphQLCancellationDoesNotHideReceiptPersistenceFailure(t *testing.T) { if err != nil { t.Fatal(err) } - st.Close() - client := historyFixtureClient{err: context.DeadlineExceeded} + // Close during the fetch, after the owner-exclusion preflight. This tests + // a cancelled provider request whose durable receipt cannot be persisted. + client := historyFixtureClient{err: context.DeadlineExceeded, beforeFetch: func() { st.Close() }} _, err = New(client, st).Sync(ctx, Options{Owner: "fixture", Repo: "repo", Numbers: []int{1}, State: "all", GraphQLHistory: true, IncludeComments: true, IncludePRMetadata: true, ReceiptOperation: "review_state"}) if !errors.Is(err, context.DeadlineExceeded) || !errors.Is(err, errAnalyticsReceipt) { t.Fatalf("missing distinguishable receipt failure: %v", err) diff --git a/internal/syncer/syncer.go b/internal/syncer/syncer.go index a538fd7c..5ee2617c 100644 --- a/internal/syncer/syncer.go +++ b/internal/syncer/syncer.go @@ -64,6 +64,7 @@ type Options struct { } type Stats struct { + OwnerExcluded int `json:"owner_excluded,omitempty"` ReviewStateOnly bool `json:"review_state_only,omitempty"` FetchMillis int64 `json:"fetch_ms,omitempty"` PersistMillis int64 `json:"persist_ms,omitempty"` @@ -149,22 +150,38 @@ func (s *Syncer) Sync(ctx context.Context, options Options) (result Stats, resul if err != nil { return Stats{}, err } - var history *gh.HistoryBatch - var repoRaw map[string]any + operation := options.ReceiptOperation + if operation == "" { + operation = "graphql_history" + } + if options.GraphQLHistory { if len(options.Numbers) == 0 || !options.IncludeComments || !options.IncludePRMetadata || options.IncludePRDetails || since != "" || options.Limit != 0 || state != "all" { return Stats{}, fmt.Errorf("--graphql-history requires --numbers, --state all, --include-comments and --with pr-metadata; since/limit/pr-details are unsupported") } - operation := options.ReceiptOperation - if operation == "" { - operation = "graphql_history" - } if operation != "graphql_history" && operation != "review_state" { return Stats{}, fmt.Errorf("unsupported GraphQL receipt operation") } if options.ReviewStateOnly && operation != "review_state" { return Stats{}, fmt.Errorf("review-state-only fetch requires review recovery operation") } + } + excludedCount := 0 + defer func() { result.OwnerExcluded = excludedCount }() + if len(options.Numbers) > 0 { + selected := uniquePositiveNumbers(options.Numbers) + options.Numbers, err = s.store.FilterExcludedNumbers(ctx, options.Owner+"/"+options.Repo, selected) + if err != nil { + return Stats{}, err + } + excludedCount = len(selected) - len(options.Numbers) + if len(options.Numbers) == 0 { + return Stats{Repository: options.Owner + "/" + options.Repo, StartedAt: started, FinishedAt: s.now().Format(time.RFC3339Nano)}, nil + } + } + var history *gh.HistoryBatch + var repoRaw map[string]any + if options.GraphQLHistory { // Fetch/validation failures happen before conversation transactions and // were previously invisible to durable run tables. Keep a receipt even // when the request is cancelled; accepted content remains untouched. @@ -309,6 +326,14 @@ func (s *Syncer) Sync(ctx context.Context, options Options) (result Stats, resul for _, row := range rows { payload := threadSyncPayload{row: row} number := intValue(row["number"]) + excluded, err := s.store.ThreadExcluded(ctx, options.Owner+"/"+options.Repo, number) + if err != nil { + return Stats{}, err + } + if excluded { + excludedCount++ + continue + } kind := issueKind(row) if history != nil { // Keep legacy REST identity stable when revisiting a previously saved diff --git a/internal/syncer/thread_exclusions_test.go b/internal/syncer/thread_exclusions_test.go new file mode 100644 index 00000000..eeebb8a0 --- /dev/null +++ b/internal/syncer/thread_exclusions_test.go @@ -0,0 +1,132 @@ +package syncer + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "path/filepath" + "sync/atomic" + "testing" + + gh "github.com/openclaw/gitcrawl/internal/github" + "github.com/openclaw/gitcrawl/internal/store" +) + +func TestOwnerExcludedSelectionNeverReachesProviderOrRetryLedger(t *testing.T) { + ctx := context.Background() + s, err := store.Open(ctx, filepath.Join(t.TempDir(), "archive.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + _, err = s.DB().Exec(`INSERT INTO repositories(id,owner,name,full_name,raw_json,updated_at) VALUES(1,'fixture','repo','fixture/repo','{}','2026-01-01T00:00:00Z'); + INSERT INTO threads(id,repo_id,github_id,number,kind,state,title,html_url,labels_json,assignees_json,raw_json,content_hash,updated_at) VALUES + (1,1,'one',10,'pull_request','open','removed','','[]','[]','{}','h','2026-01-01T00:00:00Z'), + (2,1,'two',20,'pull_request','open','keeper','','[]','[]','{}','h','2026-01-01T00:00:00Z')`) + if err != nil { + t.Fatal(err) + } + plan, err := s.PlanThreadPurge(ctx, "fixture/repo", []int{10}) + if err != nil { + t.Fatal(err) + } + if _, err = s.PurgeThreads(ctx, "fixture/repo", []int{10}, plan.PlanID); err != nil { + t.Fatal(err) + } + var calls atomic.Int32 + var crawl atomic.Bool + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls.Add(1) + if !crawl.Load() { + http.Error(w, "unexpected request", 400) + return + } + switch r.URL.Path { + case "/repos/fixture/repo": + fmt.Fprint(w, `{"id":100,"full_name":"fixture/repo","owner":{"login":"fixture"},"name":"repo"}`) + case "/repos/fixture/repo/issues": + fmt.Fprint(w, `[{"id":"one","number":10,"title":"removed","state":"open","body":"do not restore","pull_request":{}},{"id":"two","number":20,"title":"keeper","state":"open","body":"keep","pull_request":{}}]`) + case "/repos/fixture/repo/issues/20/comments", "/repos/fixture/repo/pulls/20/reviews", "/repos/fixture/repo/pulls/20/comments": + fmt.Fprint(w, `[]`) + default: + t.Errorf("excluded or unexpected fetch: %s", r.URL.Path) + http.Error(w, "unexpected", 400) + } + })) + defer server.Close() + client := gh.New(gh.Options{BaseURL: server.URL}) + for _, graphql := range []bool{false, true} { + for _, review := range []bool{false, true} { + if review && !graphql { + continue + } + operation := "graphql_history" + if review { + operation = "review_state" + } + result, err := New(client, s).Sync(ctx, Options{Owner: "fixture", Repo: "repo", Numbers: []int{10}, State: "all", GraphQLHistory: graphql, ReviewStateOnly: review, ReceiptOperation: operation, IncludeComments: true, IncludePRMetadata: true}) + if err != nil || result.ThreadsSynced != 0 || result.OwnerExcluded != 1 { + t.Fatalf("excluded selection: %+v %v", result, err) + } + } + } + if calls.Load() != 0 { + t.Fatal("owner-excluded provider access", calls.Load()) + } + crawl.Store(true) + stats, err := New(client, s).Sync(ctx, Options{Owner: "fixture", Repo: "repo", State: "all", IncludeComments: true}) + if err != nil || stats.ThreadsSynced != 1 || stats.OwnerExcluded != 1 { + t.Fatalf("normal crawl: %+v %v", stats, err) + } + var n int + if err := s.DB().QueryRow("SELECT count(*) FROM threads WHERE number=10").Scan(&n); err != nil || n != 0 { + t.Fatal("crawl restored excluded thread", err) + } + mixed := &excludedHistoryClient{Client: client} + _, err = New(mixed, s).Sync(ctx, Options{Owner: "fixture", Repo: "repo", State: "all", Numbers: []int{10, 20}, GraphQLHistory: true, IncludeComments: true, IncludePRMetadata: true}) + if err == nil || fmt.Sprint(mixed.numbers) != "[20]" { + t.Fatal("mixed GraphQL selection not filtered", mixed.numbers, err) + } + for _, table := range []string{"analytics_fetch_attempts", "analytics_retries"} { + var n int + s.DB().QueryRow("select count(*) from " + table + " where number=10").Scan(&n) + if n != 0 { + t.Fatal("owner exclusion turned into retry/success", table, n) + } + } +} + +type excludedHistoryClient struct { + *gh.Client + numbers []int +} + +func (c *excludedHistoryClient) FetchGraphQLHistory(_ context.Context, _, _ string, numbers []int, _ gh.Reporter) (gh.HistoryBatch, error) { + c.numbers = numbers + return gh.HistoryBatch{}, fmt.Errorf("fixture provider failure") +} + +// A late receipt for owner removal is neither a provider deletion nor recovery. +func TestOwnerExcludedLateReceipt(t *testing.T) { + ctx := context.Background() + s, err := store.Open(ctx, filepath.Join(t.TempDir(), "archive.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + if _, err = s.DB().Exec("INSERT INTO thread_exclusions VALUES('fixture/repo',10,'issue',1,'one','now','owner_requested','fixture')"); err != nil { + t.Fatal(err) + } + for _, result := range []string{"failed", "success"} { + err = s.RecordAnalyticsAttempt(ctx, store.AnalyticsAttempt{Repository: "FIXTURE/REPO", Number: 10, Operation: "graphql_history", StartedAt: "2026-01-01T00:00:00Z", FinishedAt: "2026-01-01T00:00:01Z", Status: result, Evidence: json.RawMessage(`{}`)}) + if err != nil { + t.Fatal(err) + } + } + var n int + if err = s.DB().QueryRow("SELECT count(*) FROM analytics_fetch_attempts").Scan(&n); err != nil || n != 0 { + t.Fatal("late receipt created work", err) + } +}