Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 15 additions & 4 deletions cmd/engram/doctor.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,8 @@ func cmdDoctor(cfg store.Config) {
func printDoctorUsage() {
fmt.Fprintln(os.Stdout, "usage: engram doctor [--json] [--project PROJECT] [--check CODE]")
fmt.Fprintln(os.Stdout, " engram doctor repair --project PROJECT --check CODE (--plan|--dry-run|--apply)")
fmt.Fprintln(os.Stdout, " engram doctor repair [--project PROJECT] --check "+diagnostic.CheckSyncMutationRequiredFields+" (--plan|--dry-run|--apply)")
fmt.Fprintln(os.Stdout, "note: --project is required for every repair check except "+diagnostic.CheckSyncMutationRequiredFields+", where it optionally scopes the quarantine to one project.")
fmt.Fprintln(os.Stdout, "checks: "+strings.Join(diagnostic.RegisteredCodes(), ", "))
}

Expand Down Expand Up @@ -127,7 +129,7 @@ func cmdDoctorRepair(cfg store.Config) {
project, _ = store.NormalizeProject(project)
project = strings.TrimSpace(project)
check = strings.TrimSpace(check)
if project == "" {
if project == "" && check != diagnostic.CheckSyncMutationRequiredFields {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
failDoctorRepair("--project is required")
return
}
Expand All @@ -150,6 +152,15 @@ func cmdDoctorRepair(cfg store.Config) {
return
}
defer s.Close()
if check == diagnostic.CheckSyncMutationRequiredFields {
report, err := s.QuarantineIrreparableSyncMutations(project, mode == diagnostic.RepairModeApply)
if err != nil {
failDoctorRepair(err.Error())
return
}
writeDoctorRepairJSON(report)
return
}

ctx := context.Background()
report, err := runDiagnostics(ctx, s, project, check)
Expand Down Expand Up @@ -200,7 +211,7 @@ func cmdDoctorRepair(cfg store.Config) {

func isSupportedDoctorRepairCheck(check string) bool {
switch check {
case diagnostic.CheckSessionProjectDirectoryMismatch, diagnostic.CheckManualSessionNameProjectMismatch:
case diagnostic.CheckSessionProjectDirectoryMismatch, diagnostic.CheckManualSessionNameProjectMismatch, diagnostic.CheckSyncMutationRequiredFields:
return true
default:
return false
Expand All @@ -213,8 +224,8 @@ func failDoctorRepair(message string) {
exitFunc(1)
}

func writeDoctorRepairJSON(plan diagnostic.RepairPlan) {
out, err := jsonMarshalIndent(plan, "", " ")
func writeDoctorRepairJSON(value any) {
out, err := jsonMarshalIndent(value, "", " ")
if err != nil {
fatal(err)
return
Expand Down
156 changes: 155 additions & 1 deletion cmd/engram/doctor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,7 @@ func TestCmdDoctorRepairValidation(t *testing.T) {
{name: "missing mode", args: []string{"engram", "doctor", "repair", "--project", "sias-app", "--check", "session_project_directory_mismatch"}, want: "exactly one of --plan, --dry-run, or --apply is required"},
{name: "multiple modes", args: []string{"engram", "doctor", "repair", "--project", "sias-app", "--check", "session_project_directory_mismatch", "--plan", "--apply"}, want: "exactly one of --plan, --dry-run, or --apply is required"},
{name: "missing project", args: []string{"engram", "doctor", "repair", "--check", "session_project_directory_mismatch", "--plan"}, want: "--project is required"},
{name: "unsupported check", args: []string{"engram", "doctor", "repair", "--project", "sias-app", "--check", "sync_mutation_required_fields", "--plan"}, want: "unsupported repair check"},
{name: "unsupported check", args: []string{"engram", "doctor", "repair", "--project", "sias-app", "--check", "not_real", "--plan"}, want: "unsupported repair check"},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
Expand Down Expand Up @@ -377,6 +377,160 @@ func TestCmdDoctorSyncMutationRequiredFieldsBlockedEnvelope(t *testing.T) {
}
}

func TestCmdDoctorRepairQuarantinesOnlyIrreparableMutations(t *testing.T) {
cfg := testConfig(t)
s, err := store.New(cfg)
if err != nil {
t.Fatalf("store.New: %v", err)
}
if err := s.Close(); err != nil {
t.Fatalf("close store: %v", err)
}
seedDoctorPendingMutation(t, cfg, "", store.SyncEntitySession, "poison", store.SyncOpUpsert, `{"id":"poison"}`)
seedDoctorPendingMutation(t, cfg, "", store.SyncEntitySession, "later", store.SyncOpDelete, `{"id":"later"}`)

withArgs(t, "engram", "doctor", "repair", "--check", "sync_mutation_required_fields", "--dry-run")
dryOut, dryErr := captureOutput(t, func() { cmdDoctor(cfg) })
if dryErr != "" {
t.Fatalf("dry-run stderr=%q", dryErr)
}
dry := decodeRepairPlan(t, dryOut)
if dry["applied"] != false || len(dry["actions"].([]any)) != 1 {
t.Fatalf("dry-run=%v", dry)
}

withArgs(t, "engram", "doctor", "repair", "--check", "sync_mutation_required_fields", "--apply")
applyOut, applyErr := captureOutput(t, func() { cmdDoctor(cfg) })
if applyErr != "" {
t.Fatalf("apply stderr=%q", applyErr)
}
applied := decodeRepairPlan(t, applyOut)
if applied["applied"] != true || len(applied["actions"].([]any)) != 1 {
t.Fatalf("apply=%v", applied)
}
db, err := sql.Open("sqlite", filepath.Join(cfg.DataDir, "engram.db"))
if err != nil {
t.Fatalf("sql.Open: %v", err)
}
defer db.Close()
var poison, later string
if err := db.QueryRow(`SELECT disposition FROM sync_mutations WHERE entity_key = 'poison'`).Scan(&poison); err != nil {
t.Fatalf("read poison: %v", err)
}
if err := db.QueryRow(`SELECT disposition FROM sync_mutations WHERE entity_key = 'later'`).Scan(&later); err != nil {
t.Fatalf("read later: %v", err)
}
if poison != store.SyncMutationDispositionQuarantined || later != store.SyncMutationDispositionPending {
t.Fatalf("dispositions poison=%q later=%q", poison, later)
}
}

func TestCmdDoctorRepairApplyUnblocksDoctorAndKeepsPendingWork(t *testing.T) {
cfg := testConfig(t)
s, err := store.New(cfg)
if err != nil {
t.Fatalf("store.New: %v", err)
}
if err := s.Close(); err != nil {
t.Fatalf("close store: %v", err)
}
seedDoctorPendingMutation(t, cfg, "engram", store.SyncEntitySession, "poison", store.SyncOpUpsert, `{"id":"poison"}`)
seedDoctorPendingMutation(t, cfg, "engram", store.SyncEntitySession, "keep", store.SyncOpUpsert, `{"id":"keep","directory":"/work/engram"}`)
// `engram` is enrolled on purpose: the repair contract this test pins is the
// cloud one, so the check must run past the cloud-sync gate instead of taking
// the local-only early return.
enrollDoctorProject(t, cfg, "engram")

runDoctor := func(stage string) map[string]any {
t.Helper()
withArgs(t, "engram", "doctor", "--json", "--project", "engram", "--check", "sync_mutation_required_fields")
stdout, stderr := captureOutput(t, func() { cmdDoctor(cfg) })
if stderr != "" {
t.Fatalf("%s stderr=%q", stage, stderr)
}
var report map[string]any
if err := json.Unmarshal([]byte(stdout), &report); err != nil {
t.Fatalf("%s doctor json invalid: %v\n%s", stage, err, stdout)
}
return report
}

if report := runDoctor("before repair"); report["status"] != "blocked" {
t.Fatalf("expected blocked doctor before repair, got %v", report)
}

withArgs(t, "engram", "doctor", "repair", "--project", "engram", "--check", "sync_mutation_required_fields", "--apply")
applyOut, applyErr := captureOutput(t, func() { cmdDoctor(cfg) })
if applyErr != "" {
t.Fatalf("apply stderr=%q", applyErr)
}
applied := decodeRepairPlan(t, applyOut)
if applied["applied"] != true || len(applied["actions"].([]any)) != 1 {
t.Fatalf("apply=%v", applied)
}

report := runDoctor("after repair")
if report["status"] == "blocked" {
t.Fatalf("doctor stayed blocked after quarantine repair: %v", report)
}
check := report["checks"].([]any)[0].(map[string]any)
if check["result"] == "blocked" || check["severity"] == "blocking" {
t.Fatalf("check stayed blocking after quarantine repair: %v", check)
}
findings := check["findings"].([]any)
if len(findings) != 1 {
t.Fatalf("expected the quarantined row to remain visible as evidence, got %v", findings)
}
finding := findings[0].(map[string]any)
if finding["severity"] != "info" || finding["reason_code"] != "sync_mutation_quarantined" || finding["requires_confirmation"] != false {
t.Fatalf("unexpected quarantined finding: %v", finding)
}
evidence := finding["evidence"].(map[string]any)
if evidence["entity_key"] != "poison" || evidence["disposition"] != store.SyncMutationDispositionQuarantined {
t.Fatalf("quarantined evidence lost mutation identity: %v", evidence)
}

reopened, err := store.New(cfg)
if err != nil {
t.Fatalf("reopen store: %v", err)
}
defer reopened.Close()
pending, err := reopened.HasPendingSyncMutationsForProject("engram")
if err != nil || !pending {
t.Fatalf("HasPendingSyncMutationsForProject=%v err=%v", pending, err)
}
for _, targetKey := range []string{store.DefaultSyncTargetKey, store.DefaultSyncTargetKey + ":engram"} {
state, err := reopened.GetSyncState(targetKey)
if err != nil {
t.Fatalf("state for %q: %v", targetKey, err)
}
if state.Lifecycle != store.SyncLifecyclePending {
t.Fatalf("quarantine repair masked pending work for %q: lifecycle=%q", targetKey, state.Lifecycle)
}
}
}

func TestPrintDoctorUsageMarksProjectOptionalOnlyForSyncMutationRepair(t *testing.T) {
withArgs(t, "engram", "doctor", "--help")
stdout, stderr := captureOutput(t, func() { cmdDoctor(testConfig(t)) })
if stderr != "" {
t.Fatalf("stderr=%q", stderr)
}
wantLines := []string{
"usage: engram doctor [--json] [--project PROJECT] [--check CODE]",
" engram doctor repair --project PROJECT --check CODE (--plan|--dry-run|--apply)",
" engram doctor repair [--project PROJECT] --check sync_mutation_required_fields (--plan|--dry-run|--apply)",
}
for _, line := range wantLines {
if !strings.Contains(stdout, line+"\n") {
t.Fatalf("usage missing line %q\n%s", line, stdout)
}
}
if !strings.Contains(stdout, "checks: ") {
t.Fatalf("usage lost the registered check list\n%s", stdout)
}
}

func TestCmdDoctorNonEnrolledPendingMutationsBlockedEnvelope(t *testing.T) {
cfg := testConfig(t)
seedDoctorSession(t, cfg, "manual-save-bootstrap", "bootstrap", "/work/bootstrap")
Expand Down
59 changes: 53 additions & 6 deletions internal/diagnostic/checks.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,8 +152,16 @@ func (c SyncMutationRequiredFieldsCheck) Run(ctx context.Context, scope Scope) (
if err != nil {
return CheckResult{}, err
}
findings := make([]Finding, 0)
blocking := make([]Finding, 0)
quarantined := make([]Finding, 0)
for _, mutation := range mutations {
// A quarantined row is an explicit, already-taken disposition: it no
// longer reaches transport, so it must not keep doctor blocked. It stays
// reported as non-blocking evidence of what was dropped from sync.
if strings.TrimSpace(mutation.Disposition) == store.SyncMutationDispositionQuarantined {
quarantined = append(quarantined, c.quarantinedFinding(mutation))
continue
}
validation := store.ValidateSyncMutationPayload(mutation.Entity, mutation.Op, mutation.Payload, mutation.EntityKey)
if validation.ReasonCode == "" {
continue
Expand All @@ -162,7 +170,7 @@ func (c SyncMutationRequiredFieldsCheck) Run(ctx context.Context, scope Scope) (
if strings.TrimSpace(scope.Project) != "" {
nextStep = "Run `engram cloud upgrade doctor --project " + scope.Project + "` and inspect the mutation payload before any manual repair."
}
findings = append(findings, Finding{
blocking = append(blocking, Finding{
CheckID: c.Code(),
Severity: SeverityBlocking,
ReasonCode: validation.ReasonCode,
Expand All @@ -173,20 +181,35 @@ func (c SyncMutationRequiredFieldsCheck) Run(ctx context.Context, scope Scope) (
RequiresConfirmation: true,
})
}
// Quarantined rows are already-taken dispositions, so they never count as
// work still pending delivery.
evidence := map[string]any{"pending_mutations_evaluated": len(mutations) - len(quarantined)}
if len(quarantined) > 0 {
evidence["quarantined_mutations"] = len(quarantined)
}
// Blocking findings lead the roll-up so the check summary always describes the
// work that still needs a decision rather than already-dispositioned evidence.
rollUp := func() []Finding { return append(append([]Finding{}, blocking...), quarantined...) }

// A non-enrolled backlog is only a fault on a device that actually uses
// cloud sync. The store journals sync mutations unconditionally, so on a
// local-only install every pending mutation belongs to a non-enrolled
// project by definition β€” the normal steady state, not something doctor
// should block on and answer with `engram cloud enroll`. This mirrors the
// autosync manager, which owns the same reason code and only evaluates it
// while cloud sync is configured and running.
// while cloud sync is configured and running. The gate is deliberately
// placed after the payload/quarantine pass so a local-only install still
// gets its quarantined evidence reported instead of silently dropped.
usesCloudSync, err := cloudSyncInUse(scope)
if err != nil {
return CheckResult{}, err
}
if !usesCloudSync {
return resultFromFindings(c.Code(), map[string]any{"pending_mutations_evaluated": len(mutations)}, findings), nil
return resultFromFindings(c.Code(), evidence, rollUp()), nil
}
// CountPendingNonEnrolledSyncMutations only counts rows whose disposition is
// still `pending`, so a quarantined row can never resurrect this blocking
// finding: the backlog it reports is genuinely undeliverable work.
nonEnrolledCounts, err := scope.Store.CountPendingNonEnrolledSyncMutations(store.DefaultSyncTargetKey)
if err != nil {
return CheckResult{}, err
Expand All @@ -197,7 +220,7 @@ func (c SyncMutationRequiredFieldsCheck) Run(ctx context.Context, scope Scope) (
if scopedProject != "" && project != scopedProject {
continue
}
findings = append(findings, Finding{
blocking = append(blocking, Finding{
CheckID: c.Code(),
Severity: SeverityBlocking,
ReasonCode: constants.ReasonNonEnrolledPendingMutations,
Expand All @@ -208,7 +231,31 @@ func (c SyncMutationRequiredFieldsCheck) Run(ctx context.Context, scope Scope) (
RequiresConfirmation: true,
})
}
return resultFromFindings(c.Code(), map[string]any{"pending_mutations_evaluated": len(mutations)}, findings), nil
return resultFromFindings(c.Code(), evidence, rollUp()), nil
}

func (c SyncMutationRequiredFieldsCheck) quarantinedFinding(mutation store.SyncMutation) Finding {
return Finding{
CheckID: c.Code(),
Severity: SeverityInfo,
ReasonCode: "sync_mutation_quarantined",
Message: "Sync mutation is quarantined and no longer blocks cloud replication.",
Why: "Quarantine keeps the irreparable journal row as durable local evidence while removing it from transport, so doctor reports it instead of staying blocked forever.",
Evidence: mustJSON(map[string]any{
"seq": mutation.Seq,
"target_key": mutation.TargetKey,
"project": mutation.Project,
"entity": mutation.Entity,
"op": mutation.Op,
"entity_key": mutation.EntityKey,
"disposition": mutation.Disposition,
"disposition_reason": mutation.DispositionReason,
"disposition_evidence": mutation.DispositionEvidence,
"disposition_at": mutation.DispositionAt,
}),
SafeNextStep: "No action required. Inspect the recorded disposition evidence if you need to know what was dropped from cloud sync.",
RequiresConfirmation: false,
}
}

func (c SQLiteLockContentionCheck) Run(ctx context.Context, scope Scope) (CheckResult, error) {
Expand Down
Loading