diff --git a/docs/mcp-scalable-workflows.md b/docs/mcp-scalable-workflows.md index 13e595ed..d4f80a74 100644 --- a/docs/mcp-scalable-workflows.md +++ b/docs/mcp-scalable-workflows.md @@ -14,7 +14,7 @@ github.search_repositories -> corpus.get_repositories github.sync_repository_context -> jobs.get -> corpus.get_repositories research.query_deepwiki github.sync_threads -> jobs.get -> corpus.rank_contribution_candidates -github.sync_thread_facets -> jobs.get -> corpus.get_threads +github.sync_thread_facets -> jobs.get -> corpus.get_thread_facets corpus.find_precedents -> workflow.find_related_work workflow.prepare_issue_set ``` @@ -180,7 +180,7 @@ Batch outputs preserve input order. Each item has one of these statuses: - `complete`: use the value; - `retryable`: retry that item after `retry_after_ms` when present; -- `unavailable`: follow `next_action` or acquire the missing facet explicitly; +- `unavailable`: follow the typed `recovery` plan or acquire the missing facet explicitly; - `failed`: fix the input or local failure before retrying. A durable job can succeed while its result is `partial`: job success means the @@ -188,7 +188,13 @@ bounded operation completed and recorded every item outcome. Poll concurrent jobs together with vectorized `jobs.get`, then retry only retryable items. Never interpret absent coverage as a zero, a passing check, or a lack of competing work. New job references carry a semantic `job:` reference, -`poll_after_ms`, and a machine-readable suggested `jobs.get` call. +`poll_after_ms`, and a typed `jobs.get` follow-up with its job ID. + +Facet synchronization completes on the same offline read plane: use +`corpus.get_thread_facets` for bounded coverage metadata and follow each +returned `resource_uri` through MCP `resources/read` for the persisted facet +observations. A missing repository, thread, or facet returns a versioned +`recovery` plan whose `then` calls are ordered and carry typed arguments. Repository and dossier absence have different recovery paths: diff --git a/docs/mcp-tool-redesign.md b/docs/mcp-tool-redesign.md index 7b969a64..38b32ca6 100644 --- a/docs/mcp-tool-redesign.md +++ b/docs/mcp-tool-redesign.md @@ -1,7 +1,7 @@ # MCP tool redesign -GitContribute targets `github.com/modelcontextprotocol/go-sdk` -`v1.7.0-pre.3` and negotiates MCP `2026-07-28`. The server continues to +GitContribute targets `github.com/modelcontextprotocol/go-sdk` `v1.7.0` and +negotiates MCP `2026-07-28`. The server continues to register generic SDK tools so the SDK owns input decoding and output-schema validation at the protocol boundary. diff --git a/docs/migrations/v2-tool-surface.md b/docs/migrations/v2-tool-surface.md deleted file mode 100644 index 2646c1a6..00000000 --- a/docs/migrations/v2-tool-surface.md +++ /dev/null @@ -1,27 +0,0 @@ -# Contribution operation migration - -This release intentionally removes the scalar and compatibility contribution -surfaces. Jobs created by removed operations are not replayable after the hard -cutover; callers must submit the canonical operation instead. - -| Removed tool or command | Replacement | -| --- | --- | -| `github.get_authenticated_identity` | `github.sync_pull_request_portfolio` with `selection: authored` | -| `github.sync_authored_pull_requests` | `github.sync_pull_request_portfolio` with `selection: authored` | -| `github.sync_pull_request_status` | `github.sync_pull_request_portfolio` with `selection: explicit` | -| `github.sync_portfolio` | `github.sync_pull_request_portfolio` | -| `github.hydrate_threads` | `github.sync_thread_facets` | -| `corpus.list_pull_request_portfolio` | `corpus.list_pull_requests` | -| `corpus.find_portfolio_overlaps` | `corpus.find_pull_request_overlaps` | -| `corpus.rank_threads` | `corpus.rank_contribution_candidates` | -| `job show` | `jobs get` | -| `job cancel` | `jobs cancel` | - -PR feedback and CI diagnostics are no longer incidental parts of generic -hydration. Use `github.sync_pull_request_feedback` for comments, reviews, -inline comments, and review-thread topology. Use `github.sync_ci_failures` for -current-head check and status observations. Both return durable jobs and retain -independent coverage; missing coverage is never a negative finding. - -The removed names are not aliases. Calls fail as unknown operations so agents -cannot silently continue using the scalar workflow. diff --git a/go.mod b/go.mod index 66b55085..0d13304d 100644 --- a/go.mod +++ b/go.mod @@ -11,7 +11,7 @@ require ( github.com/google/jsonschema-go v0.4.3 github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510 github.com/google/uuid v1.6.0 - github.com/modelcontextprotocol/go-sdk v1.7.0-pre.3 + github.com/modelcontextprotocol/go-sdk v1.7.0 github.com/pelletier/go-toml/v2 v2.4.3 github.com/pressly/goose/v3 v3.27.3 github.com/sethvargo/go-retry v0.4.0 diff --git a/go.sum b/go.sum index c88acce6..6989c296 100644 --- a/go.sum +++ b/go.sum @@ -98,8 +98,8 @@ github.com/mfridman/interpolate v0.0.2 h1:pnuTK7MQIxxFz1Gr+rjSIx9u7qVjf5VOoM/u6B github.com/mfridman/interpolate v0.0.2/go.mod h1:p+7uk6oE07mpE/Ik1b8EckO0O4ZXiGAfshKBWLUM9Xg= github.com/mitchellh/hashstructure/v2 v2.0.2 h1:vGKWl0YJqUNxE8d+h8f6NJLcCJrgbhC4NcD46KavDd4= github.com/mitchellh/hashstructure/v2 v2.0.2/go.mod h1:MG3aRVU/N29oo/V/IhBX8GR/zz4kQkprJgF2EVszyDE= -github.com/modelcontextprotocol/go-sdk v1.7.0-pre.3 h1:SEAY9IduDif4iApnZgpFkjFIdo3askSGZVbZIYyTy6I= -github.com/modelcontextprotocol/go-sdk v1.7.0-pre.3/go.mod h1:dL7u98E/zjJTGzEq+j30jQ8K2k1mb6LeAH4inEcSGts= +github.com/modelcontextprotocol/go-sdk v1.7.0 h1:yqjY2dsbKAC0LSuWZVBMrHgiG8ukXv6NRo0JiALay44= +github.com/modelcontextprotocol/go-sdk v1.7.0/go.mod h1:dL7u98E/zjJTGzEq+j30jQ8K2k1mb6LeAH4inEcSGts= github.com/muesli/cancelreader v0.2.2 h1:3I4Kt4BQjOR54NavqnDogx/MIoWBFa0StPA8ELUXHmA= github.com/muesli/cancelreader v0.2.2/go.mod h1:3XuTXfFS2VjM+HTLZY9Ak0l6eUKfijIfMUZ4EgX0QYo= github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 49f8346c..8bbcea71 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -338,12 +338,12 @@ func TestMCPReaderLocalReads(t *testing.T) { _, err = reader.Dossier(ctx, mcpcontract.RepoInput{Owner: "acme", Repo: "rocket"}) var dossierErr *mcpcontract.ToolError - if !errors.As(err, &dossierErr) || dossierErr.Code != "dossier_not_persisted" || len(dossierErr.SuggestedActions) != 1 || dossierErr.SuggestedActions[0].Tool != mcpcontract.ToolGetRepositories { + if !errors.As(err, &dossierErr) || dossierErr.Code != "dossier_not_persisted" || dossierErr.Recovery == nil || len(dossierErr.Recovery.Then) != 1 || dossierErr.Recovery.Then[0].Tool != mcpcontract.ToolGetRepositories { t.Fatalf("MCP dossier before build error = %+v", err) } _, err = reader.Dossier(ctx, mcpcontract.RepoInput{Owner: "acme", Repo: "missing"}) var repositoryErr *mcpcontract.ToolError - if !errors.As(err, &repositoryErr) || repositoryErr.Code != "repository_not_indexed" || len(repositoryErr.SuggestedActions) != 1 || repositoryErr.SuggestedActions[0].Tool != mcpcontract.ToolSyncRepositoryContext { + if !errors.As(err, &repositoryErr) || repositoryErr.Code != "repository_not_indexed" || repositoryErr.Recovery == nil || len(repositoryErr.Recovery.Then) != 1 || repositoryErr.Recovery.Then[0].Tool != mcpcontract.ToolSyncRepositoryContext { t.Fatalf("MCP dossier for missing repository error = %+v", err) } if _, err := svc.BuildRepositoryDossier(ctx, contracts.RepoRef{Owner: "acme", Repo: "rocket"}); err != nil { diff --git a/internal/app/hydration.go b/internal/app/hydration.go index 87008ea3..a650ed65 100644 --- a/internal/app/hydration.go +++ b/internal/app/hydration.go @@ -55,6 +55,8 @@ type HydratedFacet struct { // HydrateOptions controls selective thread hydration. type HydrateOptions struct { + // Kind selects the exact issue or pull request when a number is ambiguous. + Kind string // Facets lists the facets to retrieve. An empty list hydrates all facets // applicable to the thread kind. Facets []string @@ -110,7 +112,16 @@ func (s *Service) HydrateThread(ctx context.Context, repo contracts.RepoRef, num return nil, hydrateErr } - thread, err := c.GetThreadByNumber(ctx, repoProjection.ID, number) + var thread *corpus.Thread + if opts.Kind == "" { + thread, err = c.GetThreadByNumber(ctx, repoProjection.ID, number) + } else { + if opts.Kind != corpus.ThreadKindIssue && opts.Kind != corpus.ThreadKindPullRequest { + hydrateErr = fmt.Errorf("thread kind must be issue or pull_request") + return nil, hydrateErr + } + thread, err = c.GetThread(ctx, repoProjection.ID, opts.Kind, number) + } if err != nil { hydrateErr = fmt.Errorf("get thread: %w", err) return nil, hydrateErr diff --git a/internal/app/hydration_refresh.go b/internal/app/hydration_refresh.go index dcaf9382..abe91e8b 100644 --- a/internal/app/hydration_refresh.go +++ b/internal/app/hydration_refresh.go @@ -12,7 +12,7 @@ import ( // refreshHydrationThreadHeader fetches the current exact thread header before // child facets. It reuses the sync projection path so hydration cannot derive // coverage freshness from a stale or missing local header. -func (s *Service) refreshHydrationThreadHeader(ctx context.Context, repo contracts.RepoRef, number int) error { +func (s *Service) refreshHydrationThreadHeader(ctx context.Context, repo contracts.RepoRef, kind string, number int) error { ref := domain.RepoRef{Owner: repo.Owner, Repo: repo.Repo} if err := ref.Validate(); err != nil { return err @@ -20,6 +20,12 @@ func (s *Service) refreshHydrationThreadHeader(ctx context.Context, repo contrac if number <= 0 { return errors.New("thread number must be positive") } + if kind != "" && kind != "issue" && kind != "pull_request" { + return errors.New("thread kind must be issue or pull_request") + } + if kind == "" { + kind = "both" + } c, err := s.openCorpus(ctx) if err != nil { @@ -43,7 +49,7 @@ func (s *Service) refreshHydrationThreadHeader(ctx context.Context, repo contrac ctx: ctx, corpus: c, repositoryID: repository.ID, - kind: "both", + kind: kind, } _, err = syncExactThreadHeaders(ctx, reader, ref, []int{number}, newSyncRequestBudget(1), writer) return err diff --git a/internal/app/hydration_repo.go b/internal/app/hydration_repo.go index d10f3be1..416b53de 100644 --- a/internal/app/hydration_repo.go +++ b/internal/app/hydration_repo.go @@ -112,6 +112,7 @@ func (s *Service) HydrateRepository(ctx context.Context, repo contracts.RepoRef, } hr, err := s.HydrateThread(ctx, repo, t.Number, HydrateOptions{ + Kind: t.Kind, Facets: facets, MaxPages: maxPages, }) diff --git a/internal/app/mcp.go b/internal/app/mcp.go index 1ca6bcb3..04644807 100644 --- a/internal/app/mcp.go +++ b/internal/app/mcp.go @@ -239,10 +239,9 @@ func (r *MCPReader) Dossier(ctx context.Context, in mcpcontract.RepoInput) (mcpc return mcpcontract.DossierOutput{}, mcpcontract.Unavailable( "repository_not_indexed", fmt.Sprintf("Repository %s is not present in the local corpus.", ref), - mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolSyncRepositoryContext, - Reason: "Persist repository metadata before requesting local derived artifacts.", - Arguments: &mcpcontract.SuggestedActionArguments{ + mcpcontract.ToolCall{ + Tool: mcpcontract.ToolSyncRepositoryContext, + Arguments: &mcpcontract.ToolCallArguments{ Repositories: []mcpcontract.RepositoryRef{{Owner: ref.Owner, Repo: ref.Repo}}, }, }, @@ -256,10 +255,9 @@ func (r *MCPReader) Dossier(ctx context.Context, in mcpcontract.RepoInput) (mcpc return mcpcontract.DossierOutput{}, mcpcontract.Unavailable( "dossier_not_persisted", fmt.Sprintf("No persisted dossier exists for %s.", ref), - mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolGetRepositories, - Reason: "Read available repository metadata without creating local state.", - Arguments: &mcpcontract.SuggestedActionArguments{ + mcpcontract.ToolCall{ + Tool: mcpcontract.ToolGetRepositories, + Arguments: &mcpcontract.ToolCallArguments{ Repositories: []mcpcontract.RepositoryRef{{Owner: ref.Owner, Repo: ref.Repo}}, }, }, @@ -622,11 +620,14 @@ func (r *MCPReader) GetCoverage(ctx context.Context, in mcpcontract.GetCoverageI out.Status = "partial" } else if reason != "" { item.Status, item.Reason = "unavailable", reason - if reason == "not_indexed" { + item.Message = "owner/repo and optional kind/number must identify a repository or exact thread" + switch reason { + case "repository_not_indexed": + item.Message = "target is not present in the local corpus" + item.Recovery = recoveryPlan(reason, item.Message, syncRepositoryContextCall(target.Owner, target.Repo)) + case "thread_not_indexed": item.Message = "target is not present in the local corpus" - item.NextAction = "Synchronize the repository or thread explicitly, then retry this item." - } else { - item.Message = "owner/repo and optional kind/number must identify a repository or exact thread" + item.Recovery = recoveryPlan(reason, item.Message, syncThreadCall(mcpcontract.ThreadRef(target))) } out.Status = "partial" } else { @@ -661,7 +662,7 @@ func readCoverageTarget(ctx context.Context, c *corpus.Corpus, target mcpcontrac return mcpcontract.CoverageOutput{}, "", fmt.Errorf("get repository: %w", err) } if repo == nil { - return mcpcontract.CoverageOutput{}, "not_indexed", nil + return mcpcontract.CoverageOutput{}, "repository_not_indexed", nil } var threadID *int64 asOf := repo.SourceUpdatedAt @@ -671,7 +672,7 @@ func readCoverageTarget(ctx context.Context, c *corpus.Corpus, target mcpcontrac return mcpcontract.CoverageOutput{}, "", fmt.Errorf("get thread: %w", err) } if thread == nil { - return mcpcontract.CoverageOutput{}, "not_indexed", nil + return mcpcontract.CoverageOutput{}, "thread_not_indexed", nil } threadID = &thread.ID asOf = thread.SourceUpdatedAt diff --git a/internal/app/mcp_advanced_reads.go b/internal/app/mcp_advanced_reads.go index 28ac6155..5b8af942 100644 --- a/internal/app/mcp_advanced_reads.go +++ b/internal/app/mcp_advanced_reads.go @@ -58,11 +58,11 @@ func (r *MCPReader) FindClusters(ctx context.Context, in mcpcontract.FindCluster item.Value = &value case errors.Is(err, errRepositoryNotFound): item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", err.Error() - item.NextAction = "Call github.sync_repository_context for this repository." + item.Recovery = recoveryPlan("repository_not_indexed", err.Error(), syncRepositoryContextCall(target.Owner, target.Repo)) out.Status = "partial" case errors.Is(err, errThreadNotFound): item.Status, item.Reason, item.Message = "unavailable", "thread_not_indexed", err.Error() - item.NextAction = "Call github.sync_threads for this exact thread." + item.Recovery = recoveryPlan("thread_not_indexed", err.Error(), syncThreadCall(mcpcontract.ThreadRef(target))) out.Status = "partial" default: item.Status, item.Reason, item.Message = "failed", "read_failed", err.Error() @@ -199,11 +199,11 @@ func (r *MCPReader) FindNeighbors(ctx context.Context, in mcpcontract.FindNeighb item.Value = &value case errors.Is(err, errRepositoryNotFound): item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", err.Error() - item.NextAction = "Call github.sync_repository_context for this repository." + item.Recovery = recoveryPlan("repository_not_indexed", err.Error(), syncRepositoryContextCall(thread.Owner, thread.Repo)) out.Status = "partial" case errors.Is(err, errThreadNotFound): item.Status, item.Reason, item.Message = "unavailable", "thread_not_indexed", err.Error() - item.NextAction = "Call github.sync_threads for this exact thread." + item.Recovery = recoveryPlan("thread_not_indexed", err.Error(), syncThreadCall(mcpcontract.ThreadRef{Owner: thread.Owner, Repo: thread.Repo, Kind: thread.Kind, Number: thread.Number})) out.Status = "partial" default: item.Status, item.Reason, item.Message = "failed", "read_failed", err.Error() diff --git a/internal/app/mcp_issue_set.go b/internal/app/mcp_issue_set.go index 3dd9d6dc..0b85ab52 100644 --- a/internal/app/mcp_issue_set.go +++ b/internal/app/mcp_issue_set.go @@ -54,9 +54,9 @@ func (r *MCPReader) PrepareIssueSet(ctx context.Context, in mcpcontract.PrepareI out.Coverage = []mcpcontract.FacetCoverageOutput{{Facet: "threads", Status: "unknown"}} out.Gaps = append(out.Gaps, mcpcontract.IssueSetGap{ Code: "relationship_population_unknown", Facet: "threads", Status: "unknown", - Message: "repository thread coverage is not recorded, so the stored pull-request population may be incomplete", NextAction: action, + Message: "repository thread coverage is not recorded, so the stored pull-request population may be incomplete", Recovery: recoveryPlan("coverage_stale", "Complete the stored pull-request population used for related-work analysis.", action), }) - out.SuggestedActions = append(out.SuggestedActions, action) + out.RecoveryPlans = append(out.RecoveryPlans, *recoveryPlan("coverage_stale", "Complete the stored pull-request population used for related-work analysis.", action)) out.Status = "partial" } else { status := "complete" @@ -65,9 +65,9 @@ func (r *MCPReader) PrepareIssueSet(ctx context.Context, in mcpcontract.PrepareI action := repositoryPullRequestSyncAction(ref) out.Gaps = append(out.Gaps, mcpcontract.IssueSetGap{ Code: "relationship_population_incomplete", Facet: "threads", Status: "partial", - Message: "repository thread coverage is incomplete, so the stored pull-request population is not exhaustive", NextAction: action, + Message: "repository thread coverage is incomplete, so the stored pull-request population is not exhaustive", Recovery: recoveryPlan("facet_incomplete", "Complete the stored pull-request population used for related-work analysis.", action), }) - out.SuggestedActions = append(out.SuggestedActions, action) + out.RecoveryPlans = append(out.RecoveryPlans, *recoveryPlan("facet_incomplete", "Complete the stored pull-request population used for related-work analysis.", action)) out.Status = "partial" } out.Coverage = []mcpcontract.FacetCoverageOutput{{ @@ -135,16 +135,16 @@ func (r *MCPReader) PrepareIssueSet(ctx context.Context, in mcpcontract.PrepareI evaluatedAt := r.now() for i, number := range in.IssueNumbers { - key := fmt.Sprintf("%s/%s#%d", ref.Owner, ref.Repo, number) + key := threadRefKey(mcpcontract.ThreadRef{Owner: ref.Owner, Repo: ref.Repo, Kind: corpus.ThreadKindIssue, Number: number}) item := mcpcontract.BatchItem[mcpcontract.PreparedIssueEvidence]{Key: key, Status: "complete"} issue, ok := issuesByNumber[number] if !ok { - item.Status, item.Reason = "unavailable", "issue_not_indexed" + item.Status, item.Reason = "unavailable", "thread_not_indexed" item.Message = "issue is not present in the local corpus" - item.NextAction = "Call github.sync_threads for this exact issue." + item.Recovery = recoveryPlan("thread_not_indexed", item.Message, issueSyncAction(ref, number)) out.Items[i] = item out.Status = "partial" - out.SuggestedActions = append(out.SuggestedActions, issueSyncAction(ref, number)) + out.RecoveryPlans = append(out.RecoveryPlans, *recoveryPlan("thread_not_indexed", item.Message, issueSyncAction(ref, number))) continue } value, actions, partial, err := prepareOneIssue( @@ -164,7 +164,7 @@ func (r *MCPReader) PrepareIssueSet(ctx context.Context, in mcpcontract.PrepareI } item.Value = &value out.Items[i] = item - out.SuggestedActions = append(out.SuggestedActions, actions...) + out.RecoveryPlans = append(out.RecoveryPlans, actions...) if value.RelatedWorkTruncated { out.Truncated = true } @@ -208,11 +208,11 @@ func unavailableIssueSet(in mcpcontract.PrepareIssueSetInput, out mcpcontract.Pr out.Status = "partial" for i, number := range in.IssueNumbers { out.Items[i] = mcpcontract.BatchItem[mcpcontract.PreparedIssueEvidence]{ - Key: fmt.Sprintf("%s/%s#%d", in.Owner, in.Repo, number), Status: "unavailable", + Key: threadRefKey(mcpcontract.ThreadRef{Owner: in.Owner, Repo: in.Repo, Kind: corpus.ThreadKindIssue, Number: number}), Status: "unavailable", Reason: "repository_not_indexed", Message: "repository is not present in the local corpus", - NextAction: "Call github.sync_threads for this exact issue.", + Recovery: recoveryPlan("repository_not_indexed", "Synchronize the repository, then retry this exact issue.", syncRepositoryContextCall(in.Owner, in.Repo), issueSyncAction(domain.RepoRef{Owner: in.Owner, Repo: in.Repo}, number)), } - out.SuggestedActions = append(out.SuggestedActions, issueSyncAction(domain.RepoRef{Owner: in.Owner, Repo: in.Repo}, number)) + out.RecoveryPlans = append(out.RecoveryPlans, *recoveryPlan("repository_not_indexed", "Synchronize the repository, then retry this exact issue.", syncRepositoryContextCall(in.Owner, in.Repo), issueSyncAction(domain.RepoRef{Owner: in.Owner, Repo: in.Repo}, number))) } return out } @@ -240,7 +240,7 @@ func prepareOneIssue( relationshipPopulationComplete bool, relationshipPopulationFresh bool, evaluatedAt time.Time, -) (mcpcontract.PreparedIssueEvidence, []mcpcontract.SuggestedAction, bool, error) { +) (mcpcontract.PreparedIssueEvidence, []mcpcontract.RecoveryPlan, bool, error) { value := mcpcontract.PreparedIssueEvidence{ Number: issue.Number, Title: issue.Title, State: issue.State, StateReason: issue.StateReason, Labels: append([]string(nil), issue.Labels...), BodyStatus: "unknown", @@ -251,15 +251,15 @@ func prepareOneIssue( RequiresConfirmation: true, Basis: "The caller explicitly included this issue; the implementation relationship has not been validated.", }, } - actions := []mcpcontract.SuggestedAction{} + actions := []mcpcontract.RecoveryPlan{} partial := false if !relationshipPopulationFresh { action := repositoryPullRequestSyncAction(ref) value.Gaps = append(value.Gaps, mcpcontract.IssueSetGap{ Code: "relationship_coverage_stale", Facet: "threads", Status: "unknown", - Message: "repository relationship coverage predates the issue observation", NextAction: action, + Message: "repository relationship coverage predates the issue observation", Recovery: recoveryPlan("coverage_stale", "Repository relationship coverage predates the issue observation.", action), }) - actions, partial = append(actions, action), true + actions, partial = append(actions, *recoveryPlan("coverage_stale", "Repository relationship coverage predates the issue observation.", action)), true relationshipPopulationComplete = false } relationshipEvidenceComplete := issue.Body != "" @@ -272,9 +272,9 @@ func prepareOneIssue( action := issueSyncAction(ref, issue.Number) value.Gaps = append(value.Gaps, mcpcontract.IssueSetGap{ Code: "body_unknown", Facet: "body", Status: "unknown", - Message: "the corpus does not distinguish a known-empty body from a body that was not captured", NextAction: action, + Message: "the corpus does not distinguish a known-empty body from a body that was not captured", Recovery: recoveryPlan("facet_not_observed", "Fetch the exact issue header and body into the local corpus.", action), }) - actions, partial = append(actions, action), true + actions, partial = append(actions, *recoveryPlan("facet_not_observed", "Fetch the exact issue header and body into the local corpus.", action)), true } for _, facet := range []string{FacetIssueComments, FacetIssueTimeline} { @@ -288,9 +288,9 @@ func prepareOneIssue( action := issueHydrateAction(ref, issue.Number, facet) value.Gaps = append(value.Gaps, mcpcontract.IssueSetGap{ Code: "facet_missing", Facet: facet, Status: "unknown", - Message: "no coverage observation is stored for this facet", NextAction: action, + Message: "no coverage observation is stored for this facet", Recovery: recoveryPlan("facet_not_observed", "Complete the missing issue evidence facet.", action), }) - actions, partial = append(actions, action), true + actions, partial = append(actions, *recoveryPlan("facet_not_observed", "Complete the missing issue evidence facet.", action)), true continue } status := "complete" @@ -305,9 +305,9 @@ func prepareOneIssue( action := issueHydrateAction(ref, issue.Number, facet) value.Gaps = append(value.Gaps, mcpcontract.IssueSetGap{ Code: "facet_incomplete", Facet: facet, Status: "partial", - Message: "stored facet coverage is incomplete", NextAction: action, + Message: "stored facet coverage is incomplete", Recovery: recoveryPlan("facet_incomplete", "Complete the missing issue evidence facet.", action), }) - actions, partial = append(actions, action), true + actions, partial = append(actions, *recoveryPlan("facet_incomplete", "Complete the missing issue evidence facet.", action)), true } } @@ -344,9 +344,9 @@ func prepareOneIssue( action := repositoryHistorySyncAction(ref) value.Gaps = append(value.Gaps, mcpcontract.IssueSetGap{ Code: "precedent_evidence_unavailable", Facet: "precedents", Status: "unknown", - Message: "historical precedent evidence is unavailable for this issue", NextAction: action, + Message: "historical precedent evidence is unavailable for this issue", Recovery: recoveryPlan("coverage_stale", "Fetch closed issue and pull-request headers used for historical precedent analysis.", action), }) - actions, partial = append(actions, action), true + actions, partial = append(actions, *recoveryPlan("coverage_stale", "Fetch closed issue and pull-request headers used for historical precedent analysis.", action)), true } else { for _, match := range precedents.Value.Matches { if match.Kind == corpus.ThreadKindPullRequest && match.MergedAt != "" { @@ -365,7 +365,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco unknown := func(reasons ...string) mcpcontract.ContributionDisposition { return mcpcontract.ContributionDisposition{ Status: "unknown", Confidence: "low", Unknowns: reasons, - NextAction: "Complete the listed evidence gaps, then prepare this exact issue set again.", + Recovery: recoveryPlan("blocked", "Complete the listed evidence gaps, then prepare this exact issue set again."), } } issueRef := fmt.Sprintf("issue:#%d", issue.Number) @@ -375,7 +375,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco }) { return mcpcontract.ContributionDisposition{ Status: "blocked_by_repository_policy", Confidence: "high", EvidenceRefs: []string{issueRef}, - NextAction: "Do not create an implementation workspace unless a maintainer reopens or redirects the issue.", + Recovery: recoveryPlan("blocked", "Do not create an implementation workspace unless a maintainer reopens or redirects the issue."), } } var mergedClosing, openClosing, closedUnmerged []mcpcontract.IssueSetRelatedWork @@ -398,7 +398,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco if len(mergedClosing) > 0 { return mcpcontract.ContributionDisposition{ Status: "already_resolved_upstream", Confidence: "high", EvidenceRefs: relatedWorkRefs(mergedClosing), - NextAction: "Verify the released behavior before considering any follow-up contribution.", + Recovery: recoveryPlan("blocked", "Verify the released behavior before considering any follow-up contribution."), } } if len(missingMerge) > 0 { @@ -410,7 +410,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco if len(openClosing) > 0 { return mcpcontract.ContributionDisposition{ Status: "active_competing_work", Confidence: "high", EvidenceRefs: relatedWorkRefs(openClosing), - NextAction: "Coordinate with the active pull request before creating another implementation workspace.", + Recovery: recoveryPlan("blocked", "Coordinate with the active pull request before creating another implementation workspace."), } } if len(closedUnmerged) > 0 { @@ -419,7 +419,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco } return mcpcontract.ContributionDisposition{ Status: "needs_maintainer_alignment", Confidence: "medium", EvidenceRefs: relatedWorkRefs(closedUnmerged), - NextAction: "Confirm the desired semantics and acceptable approach with maintainers before coding.", + Recovery: recoveryPlan("blocked", "Confirm the desired semantics and acceptable approach with maintainers before coding."), } } if !strings.EqualFold(issue.State, "open") { @@ -427,7 +427,7 @@ func issueContributionDisposition(issue mcpcontract.PreparedIssueEvidence) mcpco } return mcpcontract.ContributionDisposition{ Status: "ready_to_investigate", Confidence: "medium", EvidenceRefs: []string{issueRef}, - NextAction: "Investigate the current behavior and contribution fit before creating an implementation workspace.", + Recovery: recoveryPlan("blocked", "Investigate the current behavior and contribution fit before creating an implementation workspace."), } } @@ -488,30 +488,30 @@ func preparedIssueSourceAsOf(value mcpcontract.PreparedIssueEvidence) string { return latest } -func issueSyncAction(ref domain.RepoRef, number int) mcpcontract.SuggestedAction { - return mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolSyncThreads, Reason: "Fetch the exact issue header and body into the local corpus.", - Arguments: &mcpcontract.SuggestedActionArguments{ +func issueSyncAction(ref domain.RepoRef, number int) mcpcontract.ToolCall { + return mcpcontract.ToolCall{ + Tool: mcpcontract.ToolSyncThreads, + Arguments: &mcpcontract.ToolCallArguments{ Selection: "threads", Threads: []mcpcontract.ThreadRef{{Owner: ref.Owner, Repo: ref.Repo, Kind: corpus.ThreadKindIssue, Number: number}}, }, } } -func issueHydrateAction(ref domain.RepoRef, number int, facet string) mcpcontract.SuggestedAction { - return mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolHydrateThreads, Reason: "Complete the missing issue evidence facet.", - Arguments: &mcpcontract.SuggestedActionArguments{ +func issueHydrateAction(ref domain.RepoRef, number int, facet string) mcpcontract.ToolCall { + return mcpcontract.ToolCall{ + Tool: mcpcontract.ToolHydrateThreads, + Arguments: &mcpcontract.ToolCallArguments{ Threads: []mcpcontract.ThreadRef{{Owner: ref.Owner, Repo: ref.Repo, Kind: corpus.ThreadKindIssue, Number: number}}, Facets: []string{facet}, }, } } -func repositoryPullRequestSyncAction(ref domain.RepoRef) mcpcontract.SuggestedAction { - return mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolSyncThreads, Reason: "Complete the stored pull-request population used for related-work analysis.", - Arguments: &mcpcontract.SuggestedActionArguments{ +func repositoryPullRequestSyncAction(ref domain.RepoRef) mcpcontract.ToolCall { + return mcpcontract.ToolCall{ + Tool: mcpcontract.ToolSyncThreads, + Arguments: &mcpcontract.ToolCallArguments{ Selection: "repositories", Repositories: []mcpcontract.RepositoryRef{{Owner: ref.Owner, Repo: ref.Repo}}, Kind: corpus.ThreadKindPullRequest, @@ -520,10 +520,10 @@ func repositoryPullRequestSyncAction(ref domain.RepoRef) mcpcontract.SuggestedAc } } -func repositoryHistorySyncAction(ref domain.RepoRef) mcpcontract.SuggestedAction { - return mcpcontract.SuggestedAction{ - Tool: mcpcontract.ToolSyncThreads, Reason: "Fetch closed issue and pull-request headers used for historical precedent analysis.", - Arguments: &mcpcontract.SuggestedActionArguments{ +func repositoryHistorySyncAction(ref domain.RepoRef) mcpcontract.ToolCall { + return mcpcontract.ToolCall{ + Tool: mcpcontract.ToolSyncThreads, + Arguments: &mcpcontract.ToolCallArguments{ Selection: "repositories", Repositories: []mcpcontract.RepositoryRef{{Owner: ref.Owner, Repo: ref.Repo}}, Kind: "both", diff --git a/internal/app/mcp_issue_set_test.go b/internal/app/mcp_issue_set_test.go index 7103dc05..20a49e20 100644 --- a/internal/app/mcp_issue_set_test.go +++ b/internal/app/mcp_issue_set_test.go @@ -86,7 +86,7 @@ func TestPrepareIssueSetComposesStoredEvidenceWithoutClaimingClosure(t *testing. if value.Linkage.Relation != "related" || !value.Linkage.RequiresConfirmation { t.Fatalf("linkage = %+v", value.Linkage) } - if len(value.Gaps) != 1 || value.Gaps[0].Facet != FacetIssueTimeline || value.Gaps[0].NextAction.Tool != mcpcontract.ToolHydrateThreads { + if len(value.Gaps) != 1 || value.Gaps[0].Facet != FacetIssueTimeline || value.Gaps[0].Recovery == nil || len(value.Gaps[0].Recovery.Then) != 1 || value.Gaps[0].Recovery.Then[0].Tool != mcpcontract.ToolHydrateThreads { t.Fatalf("gaps = %+v", value.Gaps) } detailed, err := (&MCPReader{svc}).PrepareIssueSet(ctx, mcpcontract.PrepareIssueSetInput{ @@ -123,18 +123,18 @@ func TestPrepareIssueSetPreservesUnknownAndExactRecovery(t *testing.T) { if out.Status != "partial" || len(out.Items) != 2 { t.Fatalf("result = %+v", out) } - if len(out.Gaps) != 1 || out.Gaps[0].Code != "relationship_population_unknown" || out.Gaps[0].NextAction.Tool != mcpcontract.ToolSyncThreads { + if len(out.Gaps) != 1 || out.Gaps[0].Code != "relationship_population_unknown" || out.Gaps[0].Recovery == nil || len(out.Gaps[0].Recovery.Then) != 1 || out.Gaps[0].Recovery.Then[0].Tool != mcpcontract.ToolSyncThreads { t.Fatalf("relationship gaps = %+v", out.Gaps) } if got := out.Items[0].Value; got == nil || got.BodyStatus != "unknown" || len(got.Gaps) != 3 { t.Fatalf("known issue = %+v", got) } missing := out.Items[1] - if missing.Status != "unavailable" || missing.Reason != "issue_not_indexed" || missing.NextAction != "Call github.sync_threads for this exact issue." { + if missing.Status != "unavailable" || missing.Reason != "thread_not_indexed" || missing.Recovery == nil || len(missing.Recovery.Then) != 1 || missing.Recovery.Then[0].Tool != mcpcontract.ToolSyncThreads { t.Fatalf("missing issue = %+v", missing) } - if len(out.SuggestedActions) == 0 || out.SuggestedActions[0].Tool != mcpcontract.ToolSyncThreads { - t.Fatalf("suggested actions = %+v", out.SuggestedActions) + if len(out.RecoveryPlans) == 0 || len(out.RecoveryPlans[0].Then) != 1 || out.RecoveryPlans[0].Then[0].Tool != mcpcontract.ToolSyncThreads { + t.Fatalf("recovery plans = %+v", out.RecoveryPlans) } } diff --git a/internal/app/mcp_jobs.go b/internal/app/mcp_jobs.go index 5f7ffb5d..4c4b074d 100644 --- a/internal/app/mcp_jobs.go +++ b/internal/app/mcp_jobs.go @@ -89,7 +89,7 @@ func jobResultItem(item mcpcontract.BatchItem[mcpcontract.GetJobOutput], job *co value := jobResultToMCP(job, true) item.Value = &value if value.Status == "running" { - item.NextAction = "Poll jobs.get until this job reaches a terminal state." + item.Recovery = recoveryPlan("blocked", "Poll jobs.get until this job reaches a terminal state.", mcpcontract.ToolCall{Tool: mcpcontract.ToolGetJob, Arguments: &mcpcontract.ToolCallArguments{IDs: []string{value.ID}}}) } return item } @@ -122,7 +122,7 @@ func jobResultToMCP(job *contracts.JobResult, includeDetails bool) mcpcontract.G out.Artifacts, out.FollowUp = jobArtifactsAndFollowUp(job, total) case "queued", "running": out.FollowUp = &mcpcontract.JobFollowUp{ - Tool: mcpcontract.ToolGetJob, Reason: "Poll this job until execution_state is terminal.", + Tool: mcpcontract.ToolGetJob, Arguments: &mcpcontract.ToolCallArguments{IDs: []string{job.ID}}, Reason: "Poll this job until execution_state is terminal.", } } } @@ -205,9 +205,11 @@ func jobArtifactsAndFollowUp(job *contracts.JobResult, total int) ([]mcpcontract } } value := mcpcontract.NonNegativeInt(count) + follow := followUp(tool, "", reason) + follow.Arguments = jobFollowUpArguments(job) return []mcpcontract.JobArtifactReference{{ Kind: kind, Count: &value, References: references, ReferencesTruncated: len(result.Items) > len(references), - }}, followUp(tool, "", reason) + }}, follow } switch job.Kind { case "mine_repository_fix_patterns": @@ -248,8 +250,13 @@ func jobArtifactsAndFollowUp(job *contracts.JobResult, total int) ([]mcpcontract } case "sync_repository_context": return batch("repository_batch", mcpcontract.ToolGetRepositories, "Read synchronized repository facts and coverage from the offline corpus.") - case "sync_threads", jobKindSyncThreadFacets: + case "sync_threads": return batch("thread_batch", mcpcontract.ToolGetThreads, "Read synchronized thread facts and coverage from the offline corpus.") + case jobKindSyncThreadFacets: + var request mcpcontract.HydrateThreadsInput + _ = json.Unmarshal([]byte(job.Request), &request) + refs := append([]mcpcontract.ThreadRef(nil), request.Threads...) + return facetBatchArtifact(refs, request.Facets) case jobKindSyncPullRequestPortfolio: var result struct { PullRequests []string `json:"pull_requests"` @@ -322,6 +329,53 @@ func jobArtifactsAndFollowUp(job *contracts.JobResult, total int) ([]mcpcontract return nil, nil } +func jobFollowUpArguments(job *contracts.JobResult) *mcpcontract.ToolCallArguments { + switch job.Kind { + case "sync_repository_context": + var request mcpcontract.SyncRepositoryContextInput + if json.Unmarshal([]byte(job.Request), &request) == nil { + return &mcpcontract.ToolCallArguments{Repositories: append([]mcpcontract.RepositoryRef(nil), request.Repositories...)} + } + case "sync_threads": + var request mcpcontract.SyncThreadsInput + if json.Unmarshal([]byte(job.Request), &request) == nil { + return &mcpcontract.ToolCallArguments{Selection: request.Selection, Repositories: append([]mcpcontract.RepositoryRef(nil), request.Repositories...), Threads: append([]mcpcontract.ThreadRef(nil), request.Threads...), Kind: request.Kind, State: request.State} + } + case "index_repositories": + var request mcpcontract.IndexRepositoriesInput + if json.Unmarshal([]byte(job.Request), &request) == nil { + return &mcpcontract.ToolCallArguments{Repositories: indexRepositoriesToRefs(request.Repositories)} + } + } + return nil +} + +func facetBatchArtifact(refs []mcpcontract.ThreadRef, facetNames []string) ([]mcpcontract.JobArtifactReference, *mcpcontract.JobFollowUp) { + value := mcpcontract.NonNegativeInt(len(refs)) + follow := &mcpcontract.JobFollowUp{ + Tool: mcpcontract.ToolGetThreadFacets, + Arguments: &mcpcontract.ToolCallArguments{Threads: refs, Facets: append([]string(nil), facetNames...)}, + Reason: "Read the synchronized facet coverage and canonical facet resources from the offline corpus.", + } + return []mcpcontract.JobArtifactReference{{Kind: "thread_facet_batch", Count: &value, References: threadRefKeys(refs)}}, follow +} + +func indexRepositoriesToRefs(values []mcpcontract.IndexRepositoryInput) []mcpcontract.RepositoryRef { + refs := make([]mcpcontract.RepositoryRef, 0, len(values)) + for _, value := range values { + refs = append(refs, mcpcontract.RepositoryRef{Owner: value.Owner, Repo: value.Repo}) + } + return refs +} + +func threadRefKeys(refs []mcpcontract.ThreadRef) []string { + keys := make([]string, 0, len(refs)) + for _, ref := range refs { + keys = append(keys, threadRefKey(ref)) + } + return keys +} + func decodeJobProgress(job *contracts.JobResult) (string, int, int) { phase := strings.TrimSpace(job.Progress) var counts struct { @@ -331,19 +385,5 @@ func decodeJobProgress(job *contracts.JobResult) (string, int, int) { if json.Unmarshal([]byte(job.Statistics), &counts) == nil && counts.TotalItems >= 0 && counts.CompletedItems >= 0 { return phase, counts.CompletedItems, counts.TotalItems } - // Older durable rows used percentages and key=value statistics. Keep them - // readable while exposing only the structured MCP contract. - if strings.HasSuffix(phase, "%") { - phase = job.Kind - } - for _, field := range strings.Fields(job.Statistics) { - var n int - if _, err := fmt.Sscanf(field, "completed=%d", &n); err == nil { - counts.CompletedItems = n - } - if _, err := fmt.Sscanf(field, "total=%d", &n); err == nil { - counts.TotalItems = n - } - } return phase, counts.CompletedItems, counts.TotalItems } diff --git a/internal/app/mcp_jobs_test.go b/internal/app/mcp_jobs_test.go index 05884592..9765e79e 100644 --- a/internal/app/mcp_jobs_test.go +++ b/internal/app/mcp_jobs_test.go @@ -53,6 +53,27 @@ func TestGetJobOutputHidesLegacyStatusFromModelVisibleJSON(t *testing.T) { } } +func TestRecoveryPlanUsesVersionedTypedCallsOnTheWire(t *testing.T) { + t.Parallel() + value := mcpcontract.BatchItem[struct{}]{ + Key: "acme/rocket/pull_request#7", Status: "unavailable", + Recovery: &mcpcontract.RecoveryPlan{ + Version: mcpcontract.RecoveryPlanVersion, Reason: "thread_not_indexed", Message: "sync the exact thread", + Then: []mcpcontract.ToolCall{{Tool: mcpcontract.ToolSyncThreads, Arguments: &mcpcontract.ToolCallArguments{ + Selection: "threads", Threads: []mcpcontract.ThreadRef{{Owner: "acme", Repo: "rocket", Kind: "pull_request", Number: 7}}, + }}}, + }, + } + data, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + encoded := string(data) + if strings.Contains(encoded, "next_action") || !strings.Contains(encoded, `"version":"`+mcpcontract.RecoveryPlanVersion+`"`) || !strings.Contains(encoded, `"tool":"github.sync_threads"`) || !strings.Contains(encoded, `"kind":"pull_request"`) { + t.Fatalf("recovery wire contract = %s", encoded) + } +} + func TestRemovedJobKindsDoNotExposeCompatibilityArtifacts(t *testing.T) { t.Parallel() for _, kind := range []string{"sync_portfolio", "sync_authored_pull_requests", "sync_pull_request_status", "hydrate_threads"} { diff --git a/internal/app/mcp_portfolio_relationships.go b/internal/app/mcp_portfolio_relationships.go index 5dd447a2..942ee3fe 100644 --- a/internal/app/mcp_portfolio_relationships.go +++ b/internal/app/mcp_portfolio_relationships.go @@ -27,7 +27,13 @@ func (r *MCPReader) FindPortfolioOverlaps(ctx context.Context, in mcpcontract.Fi } out := mcpcontract.FindPortfolioOverlapsOutput{Status: "complete", Items: make([]mcpcontract.BatchItem[mcpcontract.PortfolioOverlapOutput], len(in.Candidates))} candidates, candidateIndexes := collectPortfolioCandidates(in.Candidates, &out) - prIDs, missingPullRequests, err := resolvePortfolioPullRequests(ctx, c, in.PullRequests) + pullRequests := append([]mcpcontract.ThreadRef(nil), in.PullRequests...) + for i := range pullRequests { + if pullRequests[i].Kind == "" { + pullRequests[i].Kind = corpus.ThreadKindPullRequest + } + } + prIDs, missingPullRequests, err := resolvePortfolioPullRequests(ctx, c, pullRequests) if err != nil { return mcpcontract.FindPortfolioOverlapsOutput{}, err } @@ -36,7 +42,7 @@ func (r *MCPReader) FindPortfolioOverlaps(ctx context.Context, in mcpcontract.Fi } if len(prIDs) == 0 { for _, index := range candidateIndexes { - out.Items[index] = mcpcontract.BatchItem[mcpcontract.PortfolioOverlapOutput]{Key: in.Candidates[index].Kind + ":" + in.Candidates[index].Ref, Status: "unavailable", Reason: "pull_requests_not_stored", Message: "none of the requested pull requests are available in the local corpus", NextAction: "Sync the exact authored pull requests, then retry this comparison."} + out.Items[index] = mcpcontract.BatchItem[mcpcontract.PortfolioOverlapOutput]{Key: in.Candidates[index].Kind + ":" + in.Candidates[index].Ref, Status: "unavailable", Reason: "thread_not_indexed", Message: "none of the requested pull requests are available in the local corpus", Recovery: recoveryPlan("thread_not_indexed", "Sync the exact authored pull requests, then retry this comparison.", syncPullRequestCalls(pullRequests)...)} } out.Status = "partial" return out, nil @@ -51,10 +57,10 @@ func (r *MCPReader) FindPortfolioOverlaps(ctx context.Context, in mcpcontract.Fi batch := mcpcontract.BatchItem[mcpcontract.PortfolioOverlapOutput]{Key: result.Candidate.Kind + ":" + result.Candidate.Ref, Status: "complete", Value: &value} if missingPullRequests { out.Status = "partial" - batch.Status, batch.Reason, batch.NextAction = "retryable", "comparison_set_incomplete", "Sync the missing pull requests, then retry this comparison." + batch.Status, batch.Reason, batch.Recovery = "retryable", "thread_not_indexed", recoveryPlan("thread_not_indexed", "Sync the missing pull requests, then retry this comparison.", syncPullRequestCalls(pullRequests)...) } else if result.Status == "unknown" { out.Status = "partial" - batch.Status, batch.Reason, batch.NextAction = "unavailable", "coverage_missing", "Sync pull-request status and record candidate overlap signals before retrying." + batch.Status, batch.Reason, batch.Recovery = "unavailable", "candidate_signal_unavailable", recoveryPlan("candidate_signal_unavailable", "Sync pull-request status and record candidate overlap signals before retrying.", syncPullRequestCalls(pullRequests)...) } out.Items[i] = batch } @@ -89,11 +95,11 @@ func resolvePortfolioPullRequests(ctx context.Context, c *corpus.Corpus, refs [] missing = true continue } - thread, err := c.GetThreadByNumber(ctx, repo.ID, ref.Number) + thread, err := c.GetThread(ctx, repo.ID, ref.Kind, ref.Number) if err != nil { return nil, false, err } - if thread == nil || thread.Kind != corpus.ThreadKindPullRequest { + if thread == nil { missing = true continue } @@ -159,11 +165,14 @@ func resolveStoredPullRequest(ctx context.Context, c *corpus.Corpus, ref mcpcont if repo == nil { return nil, fmt.Errorf("repository %s/%s is not stored", ref.Owner, ref.Repo) } - thread, err := c.GetThreadByNumber(ctx, repo.ID, ref.Number) + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } + thread, err := c.GetThread(ctx, repo.ID, ref.Kind, ref.Number) if err != nil { return nil, err } - if thread == nil || thread.Kind != corpus.ThreadKindPullRequest { + if thread == nil { return nil, fmt.Errorf("pull request %s/%s#%d is not stored", ref.Owner, ref.Repo, ref.Number) } return thread, nil diff --git a/internal/app/mcp_portfolio_relationships_test.go b/internal/app/mcp_portfolio_relationships_test.go index 774194b1..5dc5ce43 100644 --- a/internal/app/mcp_portfolio_relationships_test.go +++ b/internal/app/mcp_portfolio_relationships_test.go @@ -27,7 +27,7 @@ func TestFindPortfolioOverlapsIsolatesInvalidCandidatesAndMissingPullRequests(t if err != nil { t.Fatal(err) } - if out.Status != "partial" || len(out.Items) != 2 || out.Items[0].Status != "failed" || out.Items[1].Status != "retryable" || out.Items[1].Reason != "comparison_set_incomplete" { + if out.Status != "partial" || len(out.Items) != 2 || out.Items[0].Status != "failed" || out.Items[1].Status != "retryable" || out.Items[1].Reason != "thread_not_indexed" { t.Fatalf("overlap batch = %+v", out) } } diff --git a/internal/app/mcp_portfolio_sync.go b/internal/app/mcp_portfolio_sync.go index d941ddc1..0d12b32c 100644 --- a/internal/app/mcp_portfolio_sync.go +++ b/internal/app/mcp_portfolio_sync.go @@ -68,7 +68,10 @@ func (r *MCPReader) SyncPortfolio(ctx context.Context, in mcpcontract.SyncPortfo if in.Selection == "explicit" { references := make([]string, len(in.PullRequests)) for i, ref := range in.PullRequests { - references[i] = fmt.Sprintf("%s/%s#%d", ref.Owner, ref.Repo, ref.Number) + if ref.Kind == "" { + ref.Kind = "pull_request" + } + references[i] = threadRefKey(ref) } refreshed := 0 failures := make([]pullRequestStatusFailure, 0) diff --git a/internal/app/mcp_pr_health.go b/internal/app/mcp_pr_health.go index e2627301..16195477 100644 --- a/internal/app/mcp_pr_health.go +++ b/internal/app/mcp_pr_health.go @@ -61,8 +61,11 @@ func (s *Service) syncPullRequestStatusBatch(ctx context.Context, in pullRequest } else { status = "partial" ref := in.PullRequests[index] + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } failure := pullRequestStatusFailure{ - Reference: fmt.Sprintf("%s/%s#%d", ref.Owner, ref.Repo, ref.Number), + Reference: threadRefKey(ref), Status: stringMapValue(item, "status"), Reason: stringMapValue(item, "reason"), Message: stringMapValue(item, "message"), @@ -101,11 +104,14 @@ func stringMapValue(item map[string]any, key string) string { } func (s *Service) syncOnePullRequestStatus(ctx context.Context, ref mcpcontract.ThreadRef, maxPages int) map[string]any { - key := fmt.Sprintf("%s/%s#%d", ref.Owner, ref.Repo, ref.Number) if ref.Number <= 0 { - return map[string]any{"key": key, "status": "failed", "reason": "invalid_reference", "message": "pull request number must be positive"} + return map[string]any{"key": threadRefKey(ref), "status": "failed", "reason": "invalid_reference", "message": "pull request number must be positive"} + } + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest } - hydrated, err := s.HydrateThread(ctx, contracts.RepoRef{Owner: ref.Owner, Repo: ref.Repo}, ref.Number, HydrateOptions{Facets: []string{FacetPRDetails, FacetPRReviews}, MaxPages: maxPages}) + key := threadRefKey(ref) + hydrated, err := s.HydrateThread(ctx, contracts.RepoRef{Owner: ref.Owner, Repo: ref.Repo}, ref.Number, HydrateOptions{Kind: ref.Kind, Facets: []string{FacetPRDetails, FacetPRReviews}, MaxPages: maxPages}) if err != nil { status, reason, message, retry := githubBatchError(err) return map[string]any{"key": key, "status": status, "reason": reason, "message": message, "retry_after_ms": retry} @@ -116,7 +122,7 @@ func (s *Service) syncOnePullRequestStatus(ctx context.Context, ref mcpcontract. } statusReader, ok := reader.(github.PullRequestStatusReader) if !ok { - return map[string]any{"key": key, "status": "unavailable", "reason": "status_adapter_unavailable", "next_action": "Configure a GitHub reader with pull-request status support."} + return map[string]any{"key": key, "status": "unavailable", "reason": "blocked", "message": "Configure a GitHub reader with pull-request status support.", "recovery": recoveryPlan("blocked", "Configure a GitHub reader with pull-request status support.")} } baselines, err := s.pullRequestHealthBaselines(ctx, ref) if err != nil { @@ -141,7 +147,7 @@ func (s *Service) syncOnePullRequestStatus(ctx context.Context, ref mcpcontract. item := map[string]any{"key": key, "status": itemStatus, "facets": facets, "head_sha": remote.HeadSHA} if itemStatus == "retryable" { item["reason"] = "facet_incomplete" - item["next_action"] = "Retry github.sync_pull_request_portfolio in explicit mode for this pull request." + item["recovery"] = recoveryPlan("facet_incomplete", "Retry this pull request in explicit mode to complete its facets.", syncPullRequestCalls([]mcpcontract.ThreadRef{ref})...) } return item } @@ -158,7 +164,7 @@ func (s *Service) persistPullRequestHealth(ctx context.Context, ref mcpcontract. } return nil, err } - thread, err := c.GetThreadByNumber(ctx, repo.ID, ref.Number) + thread, err := c.GetThread(ctx, repo.ID, ref.Kind, ref.Number) if err != nil || thread == nil { if err == nil { err = errors.New("pull request is not stored") @@ -177,9 +183,9 @@ func (s *Service) persistPullRequestHealth(ctx context.Context, ref mcpcontract. {name: FacetPRClosingIssues, value: remote.ClosingIssues.Items, coverage: remote.ClosingIssues.Coverage}, {name: FacetPRFiles, value: remote.Files.Items, coverage: remote.Files.Coverage}, } - results := hydratedHealthResults(hydrated) + results := hydratedHealthResults(hydrated, mcpcontract.ThreadRef{Owner: repo.Owner, Repo: repo.Name, Kind: thread.Kind, Number: thread.Number}) for _, target := range targets { - result, err := persistOneHealthFacet(ctx, c, *repo, *thread, sourceUpdatedAt, target, baselines[target.name], remote) + result, err := persistOneHealthFacet(ctx, c, *repo, *thread, ref, sourceUpdatedAt, target, baselines[target.name], remote) if err != nil { return nil, err } @@ -188,20 +194,20 @@ func (s *Service) persistPullRequestHealth(ctx context.Context, ref mcpcontract. return results, nil } -func hydratedHealthResults(facets []HydratedFacet) []map[string]any { +func hydratedHealthResults(facets []HydratedFacet, ref mcpcontract.ThreadRef) []map[string]any { results := make([]map[string]any, 0, len(facets)) for _, facet := range facets { result := map[string]any{"facet": facet.Facet, "status": "complete", "complete": facet.Complete, "fetched": facet.Count, "pages": facet.Pages} if !facet.Complete { result["status"] = "retryable" - result["next_action"] = "Retry this pull request with a larger max_pages bound." + result["recovery"] = recoveryPlan("facet_incomplete", "Retry this pull request with a larger max_pages bound.", syncPullRequestCalls([]mcpcontract.ThreadRef{ref})...) } results = append(results, result) } return results } -func persistOneHealthFacet(ctx context.Context, c *corpus.Corpus, repo corpus.Repository, thread corpus.Thread, sourceUpdatedAt time.Time, target healthFacet, baseline int64, remote github.PullRequestStatus) (map[string]any, error) { +func persistOneHealthFacet(ctx context.Context, c *corpus.Corpus, repo corpus.Repository, thread corpus.Thread, ref mcpcontract.ThreadRef, sourceUpdatedAt time.Time, target healthFacet, baseline int64, remote github.PullRequestStatus) (map[string]any, error) { applied, err := persistHealthFacet(ctx, c, repo.ID, thread.ID, sourceUpdatedAt, target, baseline) if err != nil { return nil, err @@ -214,15 +220,15 @@ func persistOneHealthFacet(ctx context.Context, c *corpus.Corpus, repo corpus.Re result := map[string]any{"facet": target.name, "complete": target.coverage.Complete, "fetched": target.coverage.Fetched, "total": target.coverage.Total, "status": "complete"} if !applied { result["status"], result["complete"] = "retryable", false - result["next_action"] = "A concurrent refresh advanced this facet; retry for a coherent snapshot." + result["recovery"] = recoveryPlan("coverage_stale", "A concurrent refresh advanced this facet; retry for a coherent snapshot.", syncPullRequestCalls([]mcpcontract.ThreadRef{ref})...) } if target.name == FacetPRMergeState && !remote.MergeState.MergeableKnown { result["status"] = "retryable" - result["next_action"] = "Retry after GitHub finishes computing mergeability." + result["recovery"] = recoveryPlan("facet_incomplete", "Retry after GitHub finishes computing mergeability.", syncPullRequestCalls([]mcpcontract.ThreadRef{ref})...) } if !target.coverage.Complete { result["status"] = "retryable" - result["next_action"] = "Retry this pull request to complete the facet after the current cursor." + result["recovery"] = recoveryPlan("facet_incomplete", "Retry this pull request to complete the facet after the current cursor.", syncPullRequestCalls([]mcpcontract.ThreadRef{ref})...) } return result, nil } @@ -301,7 +307,10 @@ func (s *Service) pullRequestHealthBaselines(ctx context.Context, ref mcpcontrac } return nil, err } - thread, err := c.GetThreadByNumber(ctx, repo.ID, ref.Number) + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } + thread, err := c.GetThread(ctx, repo.ID, ref.Kind, ref.Number) if err != nil || thread == nil { if err == nil { err = errors.New("pull request is not stored") diff --git a/internal/app/mcp_pr_workflows.go b/internal/app/mcp_pr_workflows.go index aba5b3ef..6b6b7760 100644 --- a/internal/app/mcp_pr_workflows.go +++ b/internal/app/mcp_pr_workflows.go @@ -27,7 +27,7 @@ type pullRequestWorkflowItem struct { ResourceURI string `json:"resource_uri,omitempty"` Code string `json:"code,omitempty"` Message string `json:"message,omitempty"` - NextAction string `json:"next_action,omitempty"` + Recovery *mcpcontract.RecoveryPlan `json:"recovery,omitempty"` RetryAfterMS int `json:"retry_after_ms,omitempty"` } @@ -108,6 +108,9 @@ func (r *MCPReader) syncPullRequestFeedback(ctx context.Context, in mcpcontract. budget := github.NewRequestBudget(in.MaxRequests) out := pullRequestWorkflowResult{BatchStatus: "complete", Items: make([]pullRequestWorkflowItem, len(in.PullRequests))} for index, ref := range in.PullRequests { + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } item := pullRequestWorkflowItem{Key: pullRequestKey(ref), Status: "complete"} snapshot, readErr := feedbackReader.GetPullRequestFeedback(ctx, ref.Owner, ref.Repo, ref.Number, github.PullRequestFeedbackOptions{ Channels: in.Channels, ThreadState: in.ThreadState, MaxItemsPerChannel: in.MaxItemsPerChannel, @@ -131,7 +134,7 @@ func (r *MCPReader) syncPullRequestFeedback(ctx context.Context, in mcpcontract. item.Status = "retryable" item.Code = "feedback_coverage_incomplete" item.Message = "one or more feedback channels reached max_items_per_channel" - item.NextAction = "Retry only this pull request with a larger max_items_per_channel bound." + item.Recovery = recoveryPlan("facet_incomplete", item.Message, mcpcontract.ToolCall{Tool: mcpcontract.ToolSyncPullRequestFeedback, Arguments: &mcpcontract.ToolCallArguments{PullRequests: []mcpcontract.ThreadRef{ref}, Channels: append([]string(nil), in.Channels...), ThreadState: in.ThreadState, MaxItemsPerChannel: in.MaxItemsPerChannel * 2, MaxRequests: in.MaxRequests}}) item.HeadSHA = snapshot.HeadSHA out.BatchStatus = "partial" } else { @@ -255,6 +258,9 @@ func (r *MCPReader) syncCIFailures(ctx context.Context, in mcpcontract.SyncCIFai budget := github.NewRequestBudget(in.MaxRequests) out := pullRequestWorkflowResult{BatchStatus: "complete", Items: make([]pullRequestWorkflowItem, len(in.PullRequests))} for index, ref := range in.PullRequests { + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } item := pullRequestWorkflowItem{Key: pullRequestKey(ref), Status: "complete"} snapshot, readErr := ciReader.GetPullRequestCI(ctx, ref.Owner, ref.Repo, ref.Number, github.CIFailureOptions{ MaxRuns: in.MaxRunsPerPR, MaxJobsPerRun: in.MaxJobsPerRun, MaxLogBytes: in.MaxLogBytesPerJob, Logs: in.Logs, @@ -283,7 +289,7 @@ func (r *MCPReader) syncCIFailures(ctx context.Context, in mcpcontract.SyncCIFai item.Status = "retryable" item.Code = "ci_coverage_incomplete" item.Message = "one or more CI collections reached a configured item bound" - item.NextAction = "Retry only this pull request with larger run or job bounds." + item.Recovery = recoveryPlan("facet_incomplete", item.Message, mcpcontract.ToolCall{Tool: mcpcontract.ToolSyncCIFailures, Arguments: &mcpcontract.ToolCallArguments{PullRequests: []mcpcontract.ThreadRef{ref}, Logs: in.Logs, MaxRunsPerPR: in.MaxRunsPerPR, MaxJobsPerRun: in.MaxJobsPerRun, MaxLogBytesPerJob: in.MaxLogBytesPerJob, MaxRequests: in.MaxRequests}}) item.HeadSHA = snapshot.HeadSHA out.BatchStatus = "partial" } else { @@ -329,7 +335,10 @@ func (r *MCPReader) persistPullRequestWorkflowFacet(ctx context.Context, ref mcp } return err } - thread, err := c.GetThreadByNumber(ctx, repo.ID, ref.Number) + if ref.Kind == "" { + ref.Kind = corpus.ThreadKindPullRequest + } + thread, err := c.GetThread(ctx, repo.ID, ref.Kind, ref.Number) if err != nil || thread == nil { if err == nil { err = errors.New("pull request is not stored") @@ -353,7 +362,7 @@ func workflowFailure(ref mcpcontract.ThreadRef, err error, tool string) pullRequ if errors.Is(err, github.ErrRequestBudgetExhausted) { return pullRequestWorkflowItem{ Key: pullRequestKey(ref), Status: "retryable", Code: "request_budget_exhausted", Message: err.Error(), - NextAction: "Retry only this pull request with " + tool + " and a sufficient max_requests bound.", + Recovery: recoveryPlan("request_budget_exhausted", err.Error(), workflowRetryCall(tool, ref)), } } status, code, message, retryAfterMS := githubBatchError(err) @@ -362,13 +371,18 @@ func workflowFailure(ref mcpcontract.ThreadRef, err error, tool string) pullRequ Message: message, RetryAfterMS: retryAfterMS, } if status == "retryable" { - item.NextAction = "Retry only this pull request with " + tool + " and a sufficient max_requests bound." + item.Recovery = recoveryPlan(code, message, workflowRetryCall(tool, ref)) } return item } +func workflowRetryCall(tool string, ref mcpcontract.ThreadRef) mcpcontract.ToolCall { + args := &mcpcontract.ToolCallArguments{PullRequests: []mcpcontract.ThreadRef{ref}, MaxRequests: 1000} + return mcpcontract.ToolCall{Tool: tool, Arguments: args} +} + func pullRequestKey(ref mcpcontract.ThreadRef) string { - return fmt.Sprintf("%s/%s#%d", ref.Owner, ref.Repo, ref.Number) + return threadRefKey(ref) } func allWorkflowItemsFailed(items []pullRequestWorkflowItem) bool { diff --git a/internal/app/mcp_pr_workflows_test.go b/internal/app/mcp_pr_workflows_test.go index 06fbc10b..29df27b5 100644 --- a/internal/app/mcp_pr_workflows_test.go +++ b/internal/app/mcp_pr_workflows_test.go @@ -199,7 +199,7 @@ func TestFeedbackResourceUsesPublicChannelsAndPreservesThreadSelection(t *testin func TestWorkflowFailurePreservesRetryableGitHubClassification(t *testing.T) { ref := mcpcontract.ThreadRef{Owner: "acme", Repo: "rocket", Number: 7} item := workflowFailure(ref, &github.TransientError{Cause: errors.New("head changed")}, mcpcontract.ToolSyncCIFailures) - if item.Status != "retryable" || item.Code != "transient" || item.RetryAfterMS == 0 || item.NextAction == "" { + if item.Status != "retryable" || item.Code != "transient" || item.RetryAfterMS == 0 || item.Recovery == nil || len(item.Recovery.Then) != 1 { t.Fatalf("item = %+v", item) } } diff --git a/internal/app/mcp_repository_search.go b/internal/app/mcp_repository_search.go index 9a3523c3..bf393b5d 100644 --- a/internal/app/mcp_repository_search.go +++ b/internal/app/mcp_repository_search.go @@ -134,9 +134,9 @@ func addRepositorySearchAction(out *mcpcontract.SearchGitHubRepositoriesOutput, for _, remote := range items { repositories = append(repositories, mcpcontract.RepositoryRef{Owner: remote.Owner, Repo: remote.Name}) } - out.SuggestedActions = []mcpcontract.SuggestedAction{{ - Tool: mcpcontract.ToolSyncThreads, Reason: "Fetch open issue headers only for repositories selected from these metadata results.", - Arguments: &mcpcontract.SuggestedActionArguments{Selection: "repositories", Repositories: repositories, State: "open"}, + out.RecoveryPlans = []mcpcontract.RecoveryPlan{{ + Version: mcpcontract.RecoveryPlanVersion, Reason: "coverage_stale", Message: "Fetch open issue headers only for repositories selected from these metadata results.", + Then: []mcpcontract.ToolCall{{Tool: mcpcontract.ToolSyncThreads, Arguments: &mcpcontract.ToolCallArguments{Selection: "repositories", Repositories: repositories, State: "open"}}}, }} } diff --git a/internal/app/mcp_scalable_operations.go b/internal/app/mcp_scalable_operations.go index 6bff287a..0da80b22 100644 --- a/internal/app/mcp_scalable_operations.go +++ b/internal/app/mcp_scalable_operations.go @@ -103,6 +103,7 @@ func (s *Service) syncThreadsBatch(ctx context.Context, in mcpcontract.SyncThrea type task struct { key string ref contracts.RepoRef + kind string numbers []int inputIndexes []int maxRequests int @@ -115,11 +116,15 @@ func (s *Service) syncThreadsBatch(ctx context.Context, in mcpcontract.SyncThrea } else { grouped := make(map[string]int) for inputIndex, thread := range in.Threads { - key := thread.Owner + "/" + thread.Repo + kind := thread.Kind + if kind == "" { + kind = "both" + } + key := thread.Owner + "/" + thread.Repo + "\x00" + kind index, ok := grouped[key] if !ok { grouped[key] = len(tasks) - tasks = append(tasks, task{key: key, ref: contracts.RepoRef{Owner: thread.Owner, Repo: thread.Repo}}) + tasks = append(tasks, task{key: thread.Owner + "/" + thread.Repo + "/" + kind, kind: kind, ref: contracts.RepoRef{Owner: thread.Owner, Repo: thread.Repo}}) index = len(tasks) - 1 } tasks[index].numbers = append(tasks[index].numbers, thread.Number) @@ -205,7 +210,11 @@ func (s *Service) syncThreadsBatch(ctx context.Context, in mcpcontract.SyncThrea defer wg.Done() for index := range jobs { current := tasks[index] - opts := SyncOptions{Kind: kind, State: state, Since: since, Numbers: current.numbers, MaxItems: in.LimitPerRepository, MaxPages: maxPages, MaxRequests: current.maxRequests} + currentKind := kind + if in.Selection == "threads" { + currentKind = current.kind + } + opts := SyncOptions{Kind: currentKind, State: state, Since: since, Numbers: current.numbers, MaxItems: in.LimitPerRepository, MaxPages: maxPages, MaxRequests: current.maxRequests} if len(current.numbers) > 0 { opts.State = "all" opts.Since = time.Time{} @@ -244,7 +253,7 @@ func (s *Service) syncThreadsBatch(ctx context.Context, in mcpcontract.SyncThrea delete(item, "requests") delete(item, "updated") thread := in.Threads[inputIndex] - item["key"] = fmt.Sprintf("%s/%s#%d", thread.Owner, thread.Repo, thread.Number) + item["key"] = threadRefKey(thread) results[inputIndex] = item } } @@ -327,7 +336,7 @@ func queuedJobReference(id, kind, message string) mcpcontract.JobReference { return mcpcontract.JobReference{ ID: id, Ref: "job:" + id, Kind: kind, Status: "queued", Message: message, PollAfterMS: 1000, FollowUp: &mcpcontract.JobFollowUp{ - Tool: mcpcontract.ToolGetJob, Reason: "Poll this job ID after the suggested delay.", + Tool: mcpcontract.ToolGetJob, Arguments: &mcpcontract.ToolCallArguments{IDs: []string{id}}, Reason: "Poll this job ID after the suggested delay.", }, } } @@ -481,8 +490,8 @@ func (s *Service) hydrateThreadsBatch(ctx context.Context, in mcpcontract.Hydrat defer wg.Done() for index := range jobs { current := in.Threads[index] - key := fmt.Sprintf("%s/%s#%d", current.Owner, current.Repo, current.Number) - res, err := s.Hydrate(ctx, contracts.RepoRef{Owner: current.Owner, Repo: current.Repo}, current.Number, contracts.HydrateOptions{Facets: in.Facets, MaxPages: in.MaxPages}) + key := threadRefKey(current) + res, err := s.Hydrate(ctx, contracts.RepoRef{Owner: current.Owner, Repo: current.Repo}, current.Number, contracts.HydrateOptions{Kind: current.Kind, Facets: in.Facets, MaxPages: in.MaxPages}) if err != nil { status, reason, message, retry := githubBatchError(err) results[index] = map[string]any{"key": key, "status": status, "reason": reason, "message": message, "retry_after_ms": retry} @@ -675,7 +684,8 @@ func (r *MCPReader) DeepWiki(ctx context.Context, in mcpcontract.DeepWikiInput) } out := mcpcontract.DeepWikiOutput{Status: "complete", Provider: "deepwiki", Action: in.Action, Repositories: repositories, Question: in.Question, Result: res.Text, SourceURL: res.SourceURL, RetrievedAt: formatTime(r.now()), Provenance: "derived_external"} if !res.Available { - out.Status, out.Reason, out.NextAction = "unavailable", "not_indexed_or_unavailable", "Use GitHub metadata, stored corpus data, or explicit code acquisition instead." + out.Status, out.Reason = "unavailable", "blocked" + out.Recovery = recoveryPlan("blocked", "Use GitHub metadata, stored corpus data, or explicit code acquisition instead.") return out, nil } if len(out.Result) > maxBytes { @@ -683,9 +693,9 @@ func (r *MCPReader) DeepWiki(ctx context.Context, in mcpcontract.DeepWikiInput) out.Truncated = true out.Reason = "output_limit" if in.Action == "contents" { - out.NextAction = "Call structure, then ask a focused question about the relevant section. Increase max_output_bytes only when the focused read is still incomplete." + out.Recovery = recoveryPlan("blocked", "Call structure, then ask a focused question about the relevant section. Increase max_output_bytes only when the focused read is still incomplete.", mcpcontract.ToolCall{Tool: mcpcontract.ToolQueryDeepWiki, Arguments: &mcpcontract.ToolCallArguments{Action: "structure", Repository: in.Repository}}) } else { - out.NextAction = "Narrow the question or repository set. Increase max_output_bytes only when the focused read is still incomplete." + out.Recovery = recoveryPlan("blocked", "Narrow the question or repository set. Increase max_output_bytes only when the focused read is still incomplete.") } } return out, nil @@ -706,9 +716,9 @@ func rejectDuplicateRepositoryRefs(inputs []mcpcontract.RepositoryRef) error { func rejectDuplicateThreadRefs(inputs []mcpcontract.ThreadRef) error { seen := make(map[string]struct{}, len(inputs)) for _, input := range inputs { - key := strings.ToLower(fmt.Sprintf("%s\x00%s\x00%d", input.Owner, input.Repo, input.Number)) + key := strings.ToLower(fmt.Sprintf("%s\x00%s\x00%s\x00%d", input.Owner, input.Repo, input.Kind, input.Number)) if _, ok := seen[key]; ok { - return mcpcontract.InvalidArgument("threads", fmt.Sprintf("duplicate thread %s/%s#%d", input.Owner, input.Repo, input.Number), nil) + return mcpcontract.InvalidArgument("threads", fmt.Sprintf("duplicate thread %s/%s/%s#%d", input.Owner, input.Repo, input.Kind, input.Number), nil) } seen[key] = struct{}{} } diff --git a/internal/app/mcp_scalable_reads.go b/internal/app/mcp_scalable_reads.go index 81311145..c7c025e9 100644 --- a/internal/app/mcp_scalable_reads.go +++ b/internal/app/mcp_scalable_reads.go @@ -66,8 +66,8 @@ func (r *MCPReader) GetRepositories(ctx context.Context, in mcpcontract.GetRepos } repo := repositories[corpus.RepositoryKey{Owner: ref.Owner, Name: ref.Repo}] if repo == nil { - item.Status, item.Reason, item.Message = "unavailable", "not_indexed", "repository is not present in the local corpus" - item.NextAction = "Call github.sync_repository_context for this repository." + item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", "repository is not present in the local corpus" + item.Recovery = recoveryPlan(item.Reason, item.Message, syncRepositoryContextCall(input.Owner, input.Repo)) out.Items[i] = item out.Status = "partial" continue @@ -80,7 +80,7 @@ func (r *MCPReader) GetRepositories(ctx context.Context, in mcpcontract.GetRepos } coverage := coverageByRepository[corpus.RepositoryFacetKey{RepositoryID: repo.ID, Facet: "metadata"}] if coverage == nil { - value.Metadata = mcpcontract.RepositoryMetadataOutput{Status: "missing", NextAction: "Call github.sync_repository_context for this repository."} + value.Metadata = mcpcontract.RepositoryMetadataOutput{Status: "missing", Recovery: recoveryPlan("facet_not_observed", "repository metadata is not observed", syncRepositoryContextCall(input.Owner, input.Repo))} clearRepositoryFacts(&value) } else { status := "complete" @@ -154,7 +154,7 @@ func (r *MCPReader) GetThreads(ctx context.Context, in mcpcontract.GetThreadsInp return mcpcontract.GetThreadsOutput{}, err } for i, input := range in.Threads { - key := fmt.Sprintf("%s/%s#%d", input.Owner, input.Repo, input.Number) + key := threadRefKey(input) item := mcpcontract.BatchItem[mcpcontract.ThreadOutput]{Key: key, Status: "complete"} ref := domain.RepoRef{Owner: input.Owner, Repo: input.Repo} if err := ref.Validate(); err != nil || input.Number < 1 { @@ -166,14 +166,15 @@ func (r *MCPReader) GetThreads(ctx context.Context, in mcpcontract.GetThreadsInp repo := repositories[corpus.RepositoryKey{Owner: ref.Owner, Name: ref.Repo}] if repo == nil { item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", "repository is not present in the local corpus" + item.Recovery = recoveryPlan(item.Reason, item.Message, syncRepositoryContextCall(input.Owner, input.Repo)) out.Items[i] = item out.Status = "partial" continue } thread := threads[corpus.ThreadKey{RepositoryID: repo.ID, Kind: input.Kind, Number: input.Number}] if thread == nil { - item.Status, item.Reason, item.Message = "unavailable", "not_indexed", "thread is not present in the local corpus" - item.NextAction = "Call github.sync_threads in thread selection mode with this exact reference." + item.Status, item.Reason, item.Message = "unavailable", "thread_not_indexed", "thread is not present in the local corpus" + item.Recovery = recoveryPlan(item.Reason, item.Message, syncThreadCall(input)) out.Items[i] = item out.Status = "partial" continue @@ -222,7 +223,7 @@ func (r *MCPReader) GetJobs(ctx context.Context, in mcpcontract.GetJobsInput) (m } else { job := jobResultToMCP(ptr(jobResult(stored)), in.ResponseFormat == "detailed") if in.ResponseFormat == "concise" && (job.Status == "succeeded" || job.Status == "failed" || job.Status == "cancelled") { - item.NextAction = "Call jobs.get with response_format=detailed to read typed artifact and follow-up references." + item.Recovery = recoveryPlan("blocked", "Read the detailed typed artifact and follow-up references.", mcpcontract.ToolCall{Tool: mcpcontract.ToolGetJob, Arguments: &mcpcontract.ToolCallArguments{IDs: []string{id}, ResponseFormat: "detailed"}}) } item.Value = &job } @@ -601,8 +602,8 @@ func (r *MCPReader) RankOpportunities(ctx context.Context, in mcpcontract.RankOp item := mcpcontract.BatchItem[mcpcontract.RepositoryOpportunitySummaryOutput]{Key: key, Status: "complete"} report, err := r.contributionRadarAt(ctx, contracts.RadarOptions{Repo: contracts.RepoRef{Owner: input.Owner, Repo: input.Repo}, Limit: in.MaxResultsPerRepository}, evaluationTime) if err != nil { - item.Status, item.Reason, item.Message = "unavailable", "not_indexed", err.Error() - item.NextAction = "Sync repository metadata and open issue headers before ranking." + item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", err.Error() + item.Recovery = recoveryPlan(item.Reason, item.Message, syncRepositoryContextCall(input.Owner, input.Repo), mcpcontract.ToolCall{Tool: mcpcontract.ToolSyncThreads, Arguments: &mcpcontract.ToolCallArguments{Selection: "repositories", Repositories: []mcpcontract.RepositoryRef{{Owner: input.Owner, Repo: input.Repo}}, Kind: "issue", State: "open"}}) out.Repositories[i] = item out.Status = "partial" continue @@ -720,7 +721,7 @@ func (r *MCPReader) FindPrecedents(ctx context.Context, in mcpcontract.FindPrece if err := ctx.Err(); err != nil { return mcpcontract.FindPrecedentsOutput{}, err } - key := fmt.Sprintf("%s/%s#%d", input.Owner, input.Repo, input.Number) + key := threadRefKey(input) item := mcpcontract.BatchItem[mcpcontract.PrecedentSet]{Key: key, Status: "complete"} repoKey := precedent.RepositoryKey(refs[i].Repository) snapshot := snapshotsByRepo[repoKey] diff --git a/internal/app/mcp_scalable_test.go b/internal/app/mcp_scalable_test.go index 5ddee0f9..ea931789 100644 --- a/internal/app/mcp_scalable_test.go +++ b/internal/app/mcp_scalable_test.go @@ -201,7 +201,7 @@ func TestGetCoveragePreservesTargetOrderAndMissingItems(t *testing.T) { if out.Items[0].Key != "acme/rocket" || out.Items[0].Value == nil || out.Items[0].Value.Facets[0].Facet != "metadata" { t.Fatalf("repository coverage = %+v", out.Items[0]) } - if out.Items[1].Key != "acme/missing" || out.Items[1].Status != "unavailable" || out.Items[1].Reason != "not_indexed" { + if out.Items[1].Key != "acme/missing" || out.Items[1].Status != "unavailable" || out.Items[1].Reason != "repository_not_indexed" { t.Fatalf("missing coverage = %+v", out.Items[1]) } if out.Items[2].Value == nil || out.Items[2].Value.Kind != "issue" || out.Items[2].Value.Number != 7 || out.Items[2].Value.Facets[0].Status != "incomplete" { @@ -355,7 +355,7 @@ func assertCancelJobsOutput(t *testing.T, out mcpcontract.GetJobsOutput, queuedI if out.Items[1].Status != "unavailable" || out.Items[1].Reason != "not_found" { t.Fatalf("missing cancellation = %+v", out.Items[1]) } - if out.Items[2].Value == nil || out.Items[2].Value.Status != "running" || !out.Items[2].Value.CancellationRequested || out.Items[2].Value.RetryAfterMS != 1000 || out.Items[2].NextAction == "" { + if out.Items[2].Value == nil || out.Items[2].Value.Status != "running" || !out.Items[2].Value.CancellationRequested || out.Items[2].Value.RetryAfterMS != 1000 || out.Items[2].Recovery == nil || len(out.Items[2].Recovery.Then) != 1 { t.Fatalf("running cancellation = %+v", out.Items[2]) } if out.Items[3].Status != "unavailable" || out.Items[3].Reason != "terminal" { @@ -447,7 +447,7 @@ func TestGetJobsDetailedReturnsTypedArtifactsWithoutStoredPayloads(t *testing.T) if len(concise.Items) != 1 || concise.Items[0].Value == nil || len(concise.Items[0].Value.Artifacts) != 0 { t.Fatalf("concise jobs output should remain a bounded state summary: %+v", concise) } - if concise.Items[0].NextAction == "" { + if concise.Items[0].Recovery == nil || len(concise.Items[0].Recovery.Then) != 1 { t.Fatalf("default concise terminal result lacks detailed recovery hint: %+v", concise) } detailed, err := reader.GetJobs(ctx, mcpcontract.GetJobsInput{IDs: []string{job.ID}, ResponseFormat: "detailed"}) @@ -493,7 +493,7 @@ func TestSearchGitHubRepositoriesPersistsObservedMetadata(t *testing.T) { if out.NextPage != 3 || out.ResponseFormat != "concise" || len(out.Items) != 1 || out.Items[0].Value == nil || out.Items[0].Value.Ref != "repository:acme/rocket" || *out.Items[0].Value.Stars != 9001 { t.Fatalf("live search result = %+v, options = %+v", out, reader.options) } - if out.Items[0].Value.Watchers != nil || len(out.SuggestedActions) != 1 || out.SuggestedActions[0].Tool != mcpcontract.ToolSyncThreads { + if out.Items[0].Value.Watchers != nil || len(out.RecoveryPlans) != 1 || len(out.RecoveryPlans[0].Then) != 1 || out.RecoveryPlans[0].Then[0].Tool != mcpcontract.ToolSyncThreads { t.Fatalf("concise search context = %+v", out) } if out.Items[0].Value.DossierStatus != "missing" { @@ -680,7 +680,7 @@ func TestDeepWikiUsesBoundedDefaultAndSteersFocusedRecovery(t *testing.T) { if len(out.Result) != mcpcontract.DeepWikiDefaultOutputBytes || !out.Truncated { t.Fatalf("default DeepWiki bound = %d bytes, truncated=%v", len(out.Result), out.Truncated) } - if out.Reason != "output_limit" || !strings.Contains(out.NextAction, "Call structure") || !strings.Contains(out.NextAction, "focused question") { + if out.Reason != "output_limit" || out.Recovery == nil || len(out.Recovery.Then) != 1 || out.Recovery.Then[0].Tool != mcpcontract.ToolQueryDeepWiki { t.Fatalf("missing truncation recovery guidance: %+v", out) } } @@ -710,6 +710,9 @@ func TestScalableBatchInputsRejectDuplicatesInsteadOfDroppingOutcomes(t *testing if err := rejectDuplicateThreadRefs([]mcpcontract.ThreadRef{{Owner: "one", Repo: "repo", Number: 1}, {Owner: "one", Repo: "repo", Number: 1}}); err == nil { t.Fatal("duplicate threads were silently accepted") } + if err := rejectDuplicateThreadRefs([]mcpcontract.ThreadRef{{Owner: "one", Repo: "repo", Kind: "issue", Number: 1}, {Owner: "one", Repo: "repo", Kind: "pull_request", Number: 1}}); err != nil { + t.Fatalf("issue and pull request with the same number were conflated: %v", err) + } if err := rejectDuplicateIndexRepositoryInputs([]mcpcontract.IndexRepositoryInput{{Owner: "one", Repo: "repo", Remote: "first"}, {Owner: "one", Repo: "repo", Remote: "second"}}); err == nil { t.Fatal("conflicting repository remotes were silently accepted") } diff --git a/internal/app/mcp_stdio_e2e_test.go b/internal/app/mcp_stdio_e2e_test.go index cdc7f11a..5765f1d4 100644 --- a/internal/app/mcp_stdio_e2e_test.go +++ b/internal/app/mcp_stdio_e2e_test.go @@ -218,7 +218,7 @@ func TestMCPStdioExactThreadSyncFlow(t *testing.T) { detailed := callMCPTool[mcpcontract.GetJobsOutput](ctx, t, session, mcpcontract.ToolGetJob, map[string]any{ "ids": []string{job.ID}, "response_format": "detailed", }) - assertExactThreadJobItems(t, detailed, []string{"lab/project#8", "lab/project#7"}) + assertExactThreadJobItems(t, detailed, []string{"lab/project/issue#8", "lab/project/pull_request#7"}) } inspection, err := New(config.NewPaths(&config.Env{Home: home}), "e2e", nil) if err != nil { diff --git a/internal/app/mcp_thread_facets.go b/internal/app/mcp_thread_facets.go new file mode 100644 index 00000000..057fa39d --- /dev/null +++ b/internal/app/mcp_thread_facets.go @@ -0,0 +1,242 @@ +package app + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + + "github.com/morluto/gitcontribute/internal/corpus" + "github.com/morluto/gitcontribute/internal/domain" + "github.com/morluto/gitcontribute/internal/facets" + "github.com/morluto/gitcontribute/internal/mcpcontract" +) + +// GetThreadFacets reads bounded facet coverage and resource identities from the +// local corpus. It never decodes or fetches large payloads in the tool result. +func (r *MCPReader) GetThreadFacets(ctx context.Context, in mcpcontract.GetThreadFacetsInput) (mcpcontract.GetThreadFacetsOutput, error) { + if len(in.Threads) < 1 || len(in.Threads) > 100 { + return mcpcontract.GetThreadFacetsOutput{}, errors.New("threads must contain 1 to 100 items") + } + if len(in.Facets) < 1 || len(in.Facets) > 10 { + return mcpcontract.GetThreadFacetsOutput{}, errors.New("facets must contain 1 to 10 items") + } + if err := validateFacetNames(in.Facets); err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + c, err := r.openReadOnlyCorpus(ctx) + if err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + out := mcpcontract.GetThreadFacetsOutput{Status: "complete", Items: make([]mcpcontract.BatchItem[mcpcontract.ThreadFacetsOutput], len(in.Threads))} + repositoryKeys := make([]corpus.RepositoryKey, 0, len(in.Threads)) + for _, input := range in.Threads { + if (domain.RepoRef{Owner: input.Owner, Repo: input.Repo}).Validate() == nil && input.Number > 0 { + repositoryKeys = append(repositoryKeys, corpus.RepositoryKey{Owner: input.Owner, Name: input.Repo}) + } + } + repositories, err := c.GetRepositoriesBatch(ctx, repositoryKeys) + if err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + threadKeys := make([]corpus.ThreadKey, 0, len(in.Threads)) + for _, input := range in.Threads { + if repo := repositories[corpus.RepositoryKey{Owner: input.Owner, Name: input.Repo}]; repo != nil && input.Number > 0 { + threadKeys = append(threadKeys, corpus.ThreadKey{RepositoryID: repo.ID, Kind: input.Kind, Number: input.Number}) + } + } + threads, err := c.GetThreadsBatch(ctx, threadKeys) + if err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + threadIDs := make([]int64, 0, len(threads)) + for _, thread := range threads { + threadIDs = append(threadIDs, thread.ID) + } + coverage, err := c.ListThreadCoverageBatch(ctx, threadIDs, in.Facets) + if err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + observations, err := c.ListThreadFacetObservationsBatch(ctx, threadIDs, in.Facets, 100) + if err != nil { + return mcpcontract.GetThreadFacetsOutput{}, err + } + for i, input := range in.Threads { + item := mcpcontract.BatchItem[mcpcontract.ThreadFacetsOutput]{Key: threadRefKey(input), Status: "complete"} + ref := domain.RepoRef{Owner: input.Owner, Repo: input.Repo} + if ref.Validate() != nil || (input.Kind != corpus.ThreadKindIssue && input.Kind != corpus.ThreadKindPullRequest) || input.Number < 1 { + item.Status, item.Reason, item.Message = "failed", "blocked", "invalid thread reference" + out.Status = "partial" + out.Items[i] = item + continue + } + repo := repositories[corpus.RepositoryKey{Owner: ref.Owner, Name: ref.Repo}] + if repo == nil { + item.Status, item.Reason, item.Message = "unavailable", "repository_not_indexed", "repository is not present in the local corpus" + item.Recovery = recoveryPlan("repository_not_indexed", item.Message, syncRepositoryContextCall(input.Owner, input.Repo)) + out.Status = "partial" + out.Items[i] = item + continue + } + thread := threads[corpus.ThreadKey{RepositoryID: repo.ID, Kind: input.Kind, Number: input.Number}] + if thread == nil { + item.Status, item.Reason, item.Message = "unavailable", "thread_not_indexed", "thread is not present in the local corpus" + item.Recovery = recoveryPlan("thread_not_indexed", item.Message, syncThreadCall(input)) + out.Status = "partial" + out.Items[i] = item + continue + } + value := mcpcontract.ThreadFacetsOutput{Owner: ref.Owner, Repo: ref.Repo, Kind: thread.Kind, Number: thread.Number, Facets: make([]mcpcontract.ThreadFacetOutput, 0, len(in.Facets))} + for _, facet := range in.Facets { + key := corpus.ThreadFacetKey{ThreadID: thread.ID, Facet: facet} + entry := mcpcontract.ThreadFacetOutput{Facet: facet, Status: "not_observed", ResourceURI: threadFacetURI(ref.Owner, ref.Repo, thread.Kind, thread.Number, facet)} + if cov := coverage[key]; cov != nil { + entry.Complete, entry.SourceUpdatedAt = cov.Complete, formatTime(cov.SourceUpdatedAt) + entry.Status = "complete" + if !cov.Complete { + entry.Status = "incomplete" + } + } + if batch, ok := observations[key]; ok { + entry.ObservationCount = mcpcontract.NonNegativeInt(len(batch.Observations)) + } + switch entry.Status { + case "not_observed": + entry.Recovery = recoveryPlan("facet_not_observed", "Synchronize this facet, then read its coverage again.", syncFacetCall(input, facet)) + case "incomplete": + entry.Recovery = recoveryPlan("facet_incomplete", "Synchronize this facet again to complete its bounded observation set.", syncFacetCall(input, facet)) + } + value.Facets = append(value.Facets, entry) + } + item.Value = &value + out.Items[i] = item + } + return out, nil +} + +// ThreadFacetResource is the canonical offline payload read for one stored +// facet. Resource reads are intentionally separate from bounded tool output. +func (r *MCPReader) ThreadFacetResource(ctx context.Context, owner, repo, kind string, number int, facet string) (map[string]any, error) { + if err := validateFacetNames([]string{facet}); err != nil { + return nil, err + } + c, err := r.openReadOnlyCorpus(ctx) + if err != nil { + return nil, err + } + storedRepo, err := c.GetRepository(ctx, owner, repo) + if err != nil || storedRepo == nil { + if err == nil { + err = errors.New("repository is not stored") + } + return nil, err + } + thread, err := c.GetThread(ctx, storedRepo.ID, kind, number) + if err != nil || thread == nil { + if err == nil { + err = errors.New("thread is not stored") + } + return nil, err + } + observations, _, err := c.ListFacetObservationsBounded(ctx, storedRepo.ID, &thread.ID, facet, 1000) + if err != nil { + return nil, err + } + coverage, err := c.GetCoverage(ctx, storedRepo.ID, &thread.ID, facet) + if err != nil { + return nil, err + } + observationValues := make([]any, 0, len(observations)) + out := map[string]any{ + "schema_version": "gitcontribute.thread-facet.v1", + "owner": owner, "repo": repo, "kind": thread.Kind, "number": number, "facet": facet, + "observations": observationValues, + } + for _, observation := range observations { + var payload any + if err := json.Unmarshal([]byte(observation.Payload), &payload); err != nil { + return nil, fmt.Errorf("decode %s observation: %w", facet, err) + } + observationValues = append(observationValues, map[string]any{ + "source_updated_at": formatTime(observation.SourceUpdatedAt), + "observation_sequence": observation.ObservationSequence, + "payload": payload, + }) + } + out["observations"] = observationValues + if coverage != nil { + out["coverage"] = map[string]any{"complete": coverage.Complete, "source_updated_at": formatTime(coverage.SourceUpdatedAt)} + } + return out, nil +} + +func validateFacetNames(values []string) error { + allowed := make(map[string]struct{}, len(facets.AllNames())) + for _, name := range facets.AllNames() { + allowed[name] = struct{}{} + } + seen := make(map[string]struct{}, len(values)) + for _, value := range values { + if strings.TrimSpace(value) == "" { + return errors.New("facet names must not be blank") + } + if _, ok := allowed[value]; !ok { + return fmt.Errorf("unknown facet %q", value) + } + if _, ok := seen[value]; ok { + return fmt.Errorf("duplicate facet %q", value) + } + seen[value] = struct{}{} + } + return nil +} + +func threadFacetURI(owner, repo, kind string, number int, facet string) string { + return fmt.Sprintf("gitcontribute://thread/%s/%s/%s/%d/facet/%s", owner, repo, kind, number, facet) +} + +func recoveryPlan(reason, message string, calls ...mcpcontract.ToolCall) *mcpcontract.RecoveryPlan { + return &mcpcontract.RecoveryPlan{Version: mcpcontract.RecoveryPlanVersion, Reason: reason, Message: message, Then: append([]mcpcontract.ToolCall(nil), calls...)} +} + +func syncRepositoryContextCall(owner, repo string) mcpcontract.ToolCall { + return mcpcontract.ToolCall{Tool: mcpcontract.ToolSyncRepositoryContext, Arguments: &mcpcontract.ToolCallArguments{ + Repositories: []mcpcontract.RepositoryRef{{Owner: owner, Repo: repo}}, + }} +} + +func syncThreadCall(ref mcpcontract.ThreadRef) mcpcontract.ToolCall { + return mcpcontract.ToolCall{Tool: mcpcontract.ToolSyncThreads, Arguments: &mcpcontract.ToolCallArguments{ + Selection: "threads", Threads: []mcpcontract.ThreadRef{ref}, + }} +} + +func syncThreadFacetsCall(ref mcpcontract.ThreadRef, facetNames []string) mcpcontract.ToolCall { + return mcpcontract.ToolCall{Tool: mcpcontract.ToolHydrateThreads, Arguments: &mcpcontract.ToolCallArguments{ + Threads: []mcpcontract.ThreadRef{ref}, Facets: append([]string(nil), facetNames...), + }} +} + +func syncFacetCall(ref mcpcontract.ThreadRef, facetName string) mcpcontract.ToolCall { + switch facetName { + case facets.PRChecks, facets.PRReviewThreads, facets.PRMergeState, facets.PRMergeQueue, facets.PRClosingIssues, facets.PRFiles: + return syncPullRequestCalls([]mcpcontract.ThreadRef{ref})[0] + default: + return syncThreadFacetsCall(ref, []string{facetName}) + } +} + +func syncPullRequestCalls(refs []mcpcontract.ThreadRef) []mcpcontract.ToolCall { + return []mcpcontract.ToolCall{{Tool: mcpcontract.ToolSyncPortfolio, Arguments: &mcpcontract.ToolCallArguments{ + Selection: "explicit", PullRequests: append([]mcpcontract.ThreadRef(nil), refs...), + }}} +} + +func threadRefKey(ref mcpcontract.ThreadRef) string { + kind := ref.Kind + if kind == "" { + kind = "any" + } + return fmt.Sprintf("%s/%s/%s#%d", ref.Owner, ref.Repo, kind, ref.Number) +} diff --git a/internal/app/mcp_thread_facets_test.go b/internal/app/mcp_thread_facets_test.go new file mode 100644 index 00000000..ec9513fb --- /dev/null +++ b/internal/app/mcp_thread_facets_test.go @@ -0,0 +1,102 @@ +package app + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/morluto/gitcontribute/internal/contracts" + "github.com/morluto/gitcontribute/internal/corpus" + "github.com/morluto/gitcontribute/internal/mcpcontract" +) + +func TestGetThreadFacetsIsOfflineAndReturnsCanonicalResources(t *testing.T) { + t.Parallel() + ctx := context.Background() + svc := newSearchTestService(t) + svc.SetGitHubReader(panicRadarReader{}) + now := time.Date(2026, 7, 30, 12, 0, 0, 0, time.UTC) + repo, err := svc.corpus.UpsertRepository(ctx, corpus.Repository{Owner: "acme", Name: "rocket", SourceUpdatedAt: now}, `{}`) + if err != nil { + t.Fatal(err) + } + issue, err := svc.corpus.UpsertThread(ctx, corpus.Thread{RepositoryID: repo.ID, Kind: corpus.ThreadKindIssue, Number: 7, Title: "issue"}, `{}`) + if err != nil { + t.Fatal(err) + } + pullRequest, err := svc.corpus.UpsertThread(ctx, corpus.Thread{RepositoryID: repo.ID, Kind: corpus.ThreadKindPullRequest, Number: 7, Title: "pull request"}, `{}`) + if err != nil { + t.Fatal(err) + } + if err := svc.corpus.ApplyFacetObservationSet(ctx, repo.ID, &issue.ID, FacetIssueComments, now, []corpus.FacetObservationInput{{SourceUpdatedAt: now, Payload: `{"body":"hello"}`}}, true, 0); err != nil { + t.Fatal(err) + } + if err := svc.corpus.ApplyFacetObservationSet(ctx, repo.ID, &pullRequest.ID, FacetPRDetails, now, []corpus.FacetObservationInput{{SourceUpdatedAt: now, Payload: `{"merged":true}`}}, true, 0); err != nil { + t.Fatal(err) + } + + reader := &MCPReader{svc} + out, err := reader.GetThreadFacets(ctx, mcpcontract.GetThreadFacetsInput{ + Threads: []mcpcontract.ThreadRef{ + {Owner: "acme", Repo: "rocket", Kind: corpus.ThreadKindIssue, Number: 7}, + {Owner: "acme", Repo: "rocket", Kind: corpus.ThreadKindPullRequest, Number: 7}, + }, + Facets: []string{FacetIssueComments}, + }) + if err != nil { + t.Fatal(err) + } + if out.Status != "complete" || len(out.Items) != 2 || out.Items[0].Value == nil || out.Items[1].Value == nil { + t.Fatalf("facet batch = %+v", out) + } + if out.Items[0].Value.Kind != corpus.ThreadKindIssue || out.Items[1].Value.Kind != corpus.ThreadKindPullRequest { + t.Fatalf("kind preservation = %+v", out.Items) + } + if out.Items[0].Value.Facets[0].ObservationCount != 1 || out.Items[0].Value.Facets[0].ResourceURI != "gitcontribute://thread/acme/rocket/issue/7/facet/issue_comments" { + t.Fatalf("issue facets = %+v", out.Items[0].Value.Facets) + } + if out.Items[1].Value.Facets[0].ResourceURI != "gitcontribute://thread/acme/rocket/pull_request/7/facet/issue_comments" { + t.Fatalf("pull-request facet = %+v", out.Items[1].Value.Facets) + } + + missing, err := reader.GetThreadFacets(ctx, mcpcontract.GetThreadFacetsInput{ + Threads: []mcpcontract.ThreadRef{{Owner: "acme", Repo: "rocket", Kind: corpus.ThreadKindIssue, Number: 7}}, + Facets: []string{FacetIssueTimeline}, + }) + if err != nil { + t.Fatal(err) + } + if missing.Items[0].Value == nil || missing.Items[0].Value.Facets[0].Status != "not_observed" || missing.Items[0].Value.Facets[0].Recovery == nil || missing.Items[0].Value.Facets[0].Recovery.Reason != "facet_not_observed" || missing.Items[0].Value.Facets[0].Recovery.Then[0].Tool != mcpcontract.ToolHydrateThreads { + t.Fatalf("missing facet recovery = %+v", missing.Items[0]) + } + + resource, err := reader.ThreadFacetResource(ctx, "acme", "rocket", corpus.ThreadKindPullRequest, 7, FacetPRDetails) + if err != nil { + t.Fatal(err) + } + data, err := json.Marshal(resource) + if err != nil { + t.Fatal(err) + } + if string(data) == "" || resource["schema_version"] != "gitcontribute.thread-facet.v1" || len(resource["observations"].([]any)) != 1 { + t.Fatalf("facet resource = %s", data) + } +} + +func TestFacetJobFollowUpReadsFacetSurface(t *testing.T) { + t.Parallel() + job := &contracts.JobResult{ + Kind: jobKindSyncThreadFacets, + Status: "succeeded", + Request: `{"threads":[{"owner":"acme","repo":"rocket","kind":"pull_request","number":7}],"facets":["pr_details"]}`, + Result: `{"status":"complete","items":[]}`, + } + artifacts, follow := jobArtifactsAndFollowUp(job, 1) + if len(artifacts) != 1 || artifacts[0].Kind != "thread_facet_batch" || follow == nil || follow.Tool != mcpcontract.ToolGetThreadFacets || follow.Arguments == nil { + t.Fatalf("facet job result = artifacts:%+v follow:%+v", artifacts, follow) + } + if len(follow.Arguments.Threads) != 1 || follow.Arguments.Threads[0].Kind != "pull_request" || len(follow.Arguments.Facets) != 1 || follow.Arguments.Facets[0] != "pr_details" { + t.Fatalf("facet follow-up arguments = %+v", follow.Arguments) + } +} diff --git a/internal/app/surfaces_extra.go b/internal/app/surfaces_extra.go index 0d53cdec..c2247d03 100644 --- a/internal/app/surfaces_extra.go +++ b/internal/app/surfaces_extra.go @@ -116,10 +116,10 @@ func (s *Service) PlanArchiveSync(_ context.Context, repo contracts.RepoRef, opt // Hydrate adapts the explicit CLI hydration contract to selective hydration. func (s *Service) Hydrate(ctx context.Context, repo contracts.RepoRef, number int, opts contracts.HydrateOptions) (*contracts.HydrateResult, error) { - if err := s.refreshHydrationThreadHeader(ctx, repo, number); err != nil { + if err := s.refreshHydrationThreadHeader(ctx, repo, opts.Kind, number); err != nil { return nil, fmt.Errorf("refresh thread header: %w", err) } - result, err := s.HydrateThread(ctx, repo, number, HydrateOptions{Facets: opts.Facets, MaxPages: opts.MaxPages}) + result, err := s.HydrateThread(ctx, repo, number, HydrateOptions{Kind: opts.Kind, Facets: opts.Facets, MaxPages: opts.MaxPages}) if err != nil { return nil, err } diff --git a/internal/contracts/archive_contracts.go b/internal/contracts/archive_contracts.go index 50b56ba6..437d1eef 100644 --- a/internal/contracts/archive_contracts.go +++ b/internal/contracts/archive_contracts.go @@ -86,6 +86,7 @@ type ArchiveSyncOptions struct { // HydrateOptions selects bounded child facets for one stored thread. type HydrateOptions struct { + Kind string Facets []string MaxPages int } diff --git a/internal/mcpcontract/issue_set_contracts.go b/internal/mcpcontract/issue_set_contracts.go index 4cdaa86b..a1dd6b3f 100644 --- a/internal/mcpcontract/issue_set_contracts.go +++ b/internal/mcpcontract/issue_set_contracts.go @@ -13,11 +13,11 @@ type PrepareIssueSetInput struct { // IssueSetGap identifies evidence that is absent or incomplete and gives the // exact bounded recovery call without treating the absence as negative proof. type IssueSetGap struct { - Code string `json:"code"` - Facet string `json:"facet"` - Status string `json:"status"` - Message string `json:"message"` - NextAction SuggestedAction `json:"next_action"` + Code string `json:"code"` + Facet string `json:"facet"` + Status string `json:"status"` + Message string `json:"message"` + Recovery *RecoveryPlan `json:"recovery"` } // IssueSetRelatedWork is one corpus-supported issue, pull request, or external @@ -58,11 +58,11 @@ type IssueSetLinkageCandidate struct { // ContributionDisposition is a conservative, evidence-backed recommendation // made before an implementation workspace is created. type ContributionDisposition struct { - Status string `json:"status"` - Confidence string `json:"confidence"` - EvidenceRefs []string `json:"evidence_refs,omitempty"` - Unknowns []string `json:"unknowns,omitempty"` - NextAction string `json:"next_action"` + Status string `json:"status"` + Confidence string `json:"confidence"` + EvidenceRefs []string `json:"evidence_refs,omitempty"` + Unknowns []string `json:"unknowns,omitempty"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` } // PreparedIssueEvidence contains stored facts and bounded derived evidence for @@ -102,5 +102,5 @@ type PrepareIssueSetOutput struct { RelationshipPopulation int `json:"relationship_population"` RelationshipConsidered int `json:"relationship_considered"` Truncated bool `json:"truncated"` - SuggestedActions []SuggestedAction `json:"suggested_actions,omitempty"` + RecoveryPlans []RecoveryPlan `json:"recovery_plans,omitempty"` } diff --git a/internal/mcpcontract/operation_contracts.go b/internal/mcpcontract/operation_contracts.go index bd54acd2..639aa0fd 100644 --- a/internal/mcpcontract/operation_contracts.go +++ b/internal/mcpcontract/operation_contracts.go @@ -81,9 +81,10 @@ type JobArtifactFailure struct { // JobFollowUp points to the typed read plane for a job's durable result. type JobFollowUp struct { - Tool string `json:"tool,omitempty" jsonschema:"Outcome-oriented tool to use next"` - ResourceURI string `json:"resource_uri,omitempty" jsonschema:"MCP resource URI to read next"` - Reason string `json:"reason" jsonschema:"Why this follow-up is appropriate"` + Tool string `json:"tool,omitempty" jsonschema:"Outcome-oriented tool to use next"` + Arguments *ToolCallArguments `json:"arguments,omitempty" jsonschema:"Typed arguments for the follow-up tool"` + ResourceURI string `json:"resource_uri,omitempty" jsonschema:"MCP resource URI to read next"` + Reason string `json:"reason" jsonschema:"Why this follow-up is appropriate"` } // GetJobOutput reports durable state and structured progress for a job. Stored diff --git a/internal/mcpcontract/scalable_contracts.go b/internal/mcpcontract/scalable_contracts.go index 16b4a4e7..c97ab348 100644 --- a/internal/mcpcontract/scalable_contracts.go +++ b/internal/mcpcontract/scalable_contracts.go @@ -28,24 +28,49 @@ type SearchGitHubRepositoriesInput struct { ResponseFormat string `json:"response_format,omitempty" jsonschema:"concise or detailed"` } -// SuggestedAction describes a non-mandatory follow-up with reusable arguments. -type SuggestedAction struct { - Tool string `json:"tool"` - Reason string `json:"reason"` - Arguments *SuggestedActionArguments `json:"arguments,omitempty"` -} - -// SuggestedActionArguments is the bounded union of reusable follow-up -// selections returned by GitContribute tools. -type SuggestedActionArguments struct { - IDs []string `json:"ids,omitempty"` - Selection string `json:"selection,omitempty"` - Repositories []RepositoryRef `json:"repositories,omitempty"` - Threads []ThreadRef `json:"threads,omitempty"` - Facets []string `json:"facets,omitempty"` - Kind string `json:"kind,omitempty"` - State string `json:"state,omitempty"` -} +// RecoveryPlan is the only model-visible recovery shape. Versioning keeps the +// contract explicit while Then preserves the order in which calls are made. +type RecoveryPlan struct { + Version string `json:"version"` + Reason string `json:"reason"` + Message string `json:"message"` + Then []ToolCall `json:"then,omitempty"` +} + +// ToolCall is one typed, replayable MCP call in a recovery plan. +type ToolCall struct { + Tool string `json:"tool"` + Arguments *ToolCallArguments `json:"arguments,omitempty"` +} + +// ToolCallArguments is the bounded typed argument union used by advertised +// recovery calls. Fields are intentionally limited to canonical MCP inputs. +type ToolCallArguments struct { + IDs []string `json:"ids,omitempty"` + Selection string `json:"selection,omitempty"` + Repositories []RepositoryRef `json:"repositories,omitempty"` + Threads []ThreadRef `json:"threads,omitempty"` + PullRequests []ThreadRef `json:"pull_requests,omitempty"` + Facets []string `json:"facets,omitempty"` + Kind string `json:"kind,omitempty"` + State string `json:"state,omitempty"` + MaxPages int `json:"max_pages,omitempty"` + MaxRequests int `json:"max_requests,omitempty"` + LimitPerRepository int `json:"limit_per_repository,omitempty"` + ResponseFormat string `json:"response_format,omitempty"` + Channels []string `json:"channels,omitempty"` + ThreadState string `json:"thread_state,omitempty"` + Logs string `json:"logs,omitempty"` + MaxItemsPerChannel int `json:"max_items_per_channel,omitempty"` + MaxRunsPerPR int `json:"max_runs_per_pr,omitempty"` + MaxJobsPerRun int `json:"max_jobs_per_run,omitempty"` + MaxLogBytesPerJob int `json:"max_log_bytes_per_job,omitempty"` + Action string `json:"action,omitempty"` + Repository string `json:"repository,omitempty"` + Question string `json:"question,omitempty"` +} + +const RecoveryPlanVersion = "gitcontribute.recovery.v1" // SearchWarning explains a request-specific limitation and how to improve it. type SearchWarning struct { @@ -78,17 +103,17 @@ type RepositorySearchMatch struct { // SearchGitHubRepositoriesOutput contains search results and completeness metadata. type SearchGitHubRepositoriesOutput struct { - Status string `json:"status"` - Query string `json:"query"` - Interpretation string `json:"interpretation"` - ResponseFormat string `json:"response_format"` - Page int `json:"page"` - NextPage int `json:"next_page,omitempty"` - Total int `json:"total"` - Incomplete bool `json:"incomplete"` - Items []BatchItem[RepositorySearchMatch] `json:"items"` - Warnings []SearchWarning `json:"warnings,omitempty"` - SuggestedActions []SuggestedAction `json:"suggested_actions,omitempty"` + Status string `json:"status"` + Query string `json:"query"` + Interpretation string `json:"interpretation"` + ResponseFormat string `json:"response_format"` + Page int `json:"page"` + NextPage int `json:"next_page,omitempty"` + Total int `json:"total"` + Incomplete bool `json:"incomplete"` + Items []BatchItem[RepositorySearchMatch] `json:"items"` + Warnings []SearchWarning `json:"warnings,omitempty"` + RecoveryPlans []RecoveryPlan `json:"recovery_plans,omitempty"` } // RepositoryRef identifies one GitHub repository without implying that it has @@ -117,7 +142,7 @@ type BatchItem[T any] struct { Reason string `json:"reason,omitempty"` Message string `json:"message,omitempty"` RetryAfterMS NonNegativeInt `json:"retry_after_ms,omitempty"` - NextAction string `json:"next_action,omitempty"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` } // GetRepositoriesInput selects repositories for a bounded corpus read. @@ -127,10 +152,10 @@ type GetRepositoriesInput struct { // RepositoryMetadataOutput describes the coverage of repository metadata. type RepositoryMetadataOutput struct { - Status string `json:"status"` - ObservedAt string `json:"observed_at,omitempty"` - SourceUpdatedAt string `json:"source_updated_at,omitempty"` - NextAction string `json:"next_action,omitempty"` + Status string `json:"status"` + ObservedAt string `json:"observed_at,omitempty"` + SourceUpdatedAt string `json:"source_updated_at,omitempty"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` } // TypedRepositoryOutput contains repository facts with explicit metadata coverage. @@ -174,6 +199,40 @@ type GetThreadsOutput struct { Items []BatchItem[ThreadOutput] `json:"items"` } +// GetThreadFacetsInput selects bounded offline facet metadata for exact +// threads. Payloads larger than a compact result are read through the returned +// resource URI. +type GetThreadFacetsInput struct { + Threads []ThreadRef `json:"threads" jsonschema:"One to 100 exact thread identities"` + Facets []string `json:"facets" jsonschema:"One to 10 stored facet names"` +} + +// ThreadFacetOutput describes one stored facet and its canonical resource. +type ThreadFacetOutput struct { + Facet string `json:"facet"` + Status string `json:"status"` + Complete bool `json:"complete"` + ObservationCount NonNegativeInt `json:"observation_count"` + SourceUpdatedAt string `json:"source_updated_at,omitempty"` + ResourceURI string `json:"resource_uri"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` +} + +// ThreadFacetsOutput contains bounded facet metadata for one exact thread. +type ThreadFacetsOutput struct { + Owner string `json:"owner"` + Repo string `json:"repo"` + Kind string `json:"kind"` + Number int `json:"number"` + Facets []ThreadFacetOutput `json:"facets"` +} + +// GetThreadFacetsOutput preserves exact-thread input order and item failures. +type GetThreadFacetsOutput struct { + Status string `json:"batch_status"` + Items []BatchItem[ThreadFacetsOutput] `json:"items"` +} + // FindPrecedentsInput selects source threads for offline analogue discovery. type FindPrecedentsInput struct { Threads []ThreadRef `json:"threads" jsonschema:"One to 20 source threads"` @@ -455,18 +514,18 @@ type DeepWikiInput struct { // DeepWikiOutput labels provider prose as derived external content and reports // provider-level unavailability without persisting the response. type DeepWikiOutput struct { - Status string `json:"status"` - Provider string `json:"provider"` - Action string `json:"action"` - Repositories []string `json:"repositories"` - Question string `json:"question,omitempty"` - Result string `json:"result,omitempty"` - SourceURL string `json:"source_url,omitempty"` - RetrievedAt string `json:"retrieved_at"` - Provenance string `json:"provenance"` - Truncated bool `json:"truncated"` - Reason string `json:"reason,omitempty"` - NextAction string `json:"next_action,omitempty"` + Status string `json:"status"` + Provider string `json:"provider"` + Action string `json:"action"` + Repositories []string `json:"repositories"` + Question string `json:"question,omitempty"` + Result string `json:"result,omitempty"` + SourceURL string `json:"source_url,omitempty"` + RetrievedAt string `json:"retrieved_at"` + Provenance string `json:"provenance"` + Truncated bool `json:"truncated"` + Reason string `json:"reason,omitempty"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` } // ErrNotFound lets readers distinguish absent corpus objects from failures. diff --git a/internal/mcpcontract/tool_contracts.go b/internal/mcpcontract/tool_contracts.go index 5130dffe..41db0e42 100644 --- a/internal/mcpcontract/tool_contracts.go +++ b/internal/mcpcontract/tool_contracts.go @@ -8,12 +8,12 @@ import ( // ToolError is the stable, actionable shape for agent-correctable requests. type ToolError struct { - Code string `json:"code"` - Message string `json:"message"` - Field string `json:"path,omitempty"` - Retryable bool `json:"retryable"` - Example map[string]any `json:"example,omitempty"` - SuggestedActions []SuggestedAction `json:"suggested_actions,omitempty"` + Code string `json:"code"` + Message string `json:"message"` + Field string `json:"path,omitempty"` + Retryable bool `json:"retryable"` + Example map[string]any `json:"example,omitempty"` + Recovery *RecoveryPlan `json:"recovery,omitempty"` } func (e *ToolError) Error() string { @@ -32,12 +32,10 @@ func InvalidArgument(field, message string, example map[string]any) error { // Unavailable reports an agent-readable terminal state with explicit recovery // actions. It is used when retrying the same read cannot make the object // available without a distinct acquisition or build step. -func Unavailable(code, message string, actions ...SuggestedAction) error { +func Unavailable(code, message string, actions ...ToolCall) error { return &ToolError{ - Code: code, - Message: message, - Retryable: false, - SuggestedActions: append([]SuggestedAction(nil), actions...), + Code: code, Message: message, Retryable: false, + Recovery: &RecoveryPlan{Version: RecoveryPlanVersion, Reason: code, Message: message, Then: append([]ToolCall(nil), actions...)}, } } @@ -48,6 +46,7 @@ const ( ToolSearchCode = "corpus.search_code" ToolGetRepositories = "corpus.get_repositories" ToolGetThreads = "corpus.get_threads" + ToolGetThreadFacets = "corpus.get_thread_facets" ToolRankThreads = "corpus.rank_contribution_candidates" ToolFindPrecedents = "corpus.find_precedents" ToolPrepareIssueSet = "workflow.prepare_issue_set" diff --git a/internal/mcpserver/capabilities_test.go b/internal/mcpserver/capabilities_test.go index ebc1d1e5..0fed5b5a 100644 --- a/internal/mcpserver/capabilities_test.go +++ b/internal/mcpserver/capabilities_test.go @@ -39,6 +39,9 @@ func (*fakeOptionalCapabilities) GetRepositories(_ context.Context, in mcpcontra func (*fakeOptionalCapabilities) GetThreads(context.Context, mcpcontract.GetThreadsInput) (mcpcontract.GetThreadsOutput, error) { return mcpcontract.GetThreadsOutput{Status: "complete"}, nil } +func (*fakeOptionalCapabilities) GetThreadFacets(context.Context, mcpcontract.GetThreadFacetsInput) (mcpcontract.GetThreadFacetsOutput, error) { + return mcpcontract.GetThreadFacetsOutput{Status: "complete"}, nil +} func (f *fakeOptionalCapabilities) RankOpportunities(context.Context, mcpcontract.RankOpportunitiesInput) (mcpcontract.RankOpportunitiesOutput, error) { score := 87 if f.base.radarScore != 0 { @@ -130,6 +133,8 @@ type completeTestReader struct { mcpcontract.Reader NeighborReader ScalableReader + ThreadFacetReader + threadFacetResourceReader IssueSetReader PortfolioReader GitHubOperator @@ -154,7 +159,7 @@ type completeTestReader struct { func completeFakeReader(base *fakeReader) mcpcontract.Reader { optional := &fakeOptionalCapabilities{base: base} return completeTestReader{ - Reader: base, NeighborReader: optional, ScalableReader: optional, IssueSetReader: optional, + Reader: base, NeighborReader: optional, ScalableReader: optional, ThreadFacetReader: optional, threadFacetResourceReader: base, IssueSetReader: optional, PortfolioReader: optional, GitHubOperator: optional, PullRequestFeedbackOperator: optional, CIFailureOperator: optional, FixPatternOperator: optional, FixPatternReader: base, CodeIndexer: optional, MergeConflictReader: optional, ResearchReader: optional, diff --git a/internal/mcpserver/catalog_test.go b/internal/mcpserver/catalog_test.go index 57a45057..71c8f7ae 100644 --- a/internal/mcpserver/catalog_test.go +++ b/internal/mcpserver/catalog_test.go @@ -11,6 +11,7 @@ import ( "unicode" "github.com/modelcontextprotocol/go-sdk/mcp" + "github.com/morluto/gitcontribute/internal/facets" "github.com/morluto/gitcontribute/internal/mcpcontract" ) @@ -217,6 +218,9 @@ func TestToolSchemasExposeMachineReadableContracts(t *testing.T) { assertSchemaValue(t, tools[mcpcontract.ToolSearchThreads].InputSchema, []string{"properties", "kind", "enum"}, []any{"issue", "pull_request"}) assertSchemaValue(t, tools[mcpcontract.ToolSearchThreads].InputSchema, []string{"properties", "limit", "default"}, float64(20)) assertSchemaValue(t, tools[mcpcontract.ToolSearchThreads].InputSchema, []string{"properties", "limit", "maximum"}, float64(100)) + assertSchemaValue(t, tools[mcpcontract.ToolGetThreadFacets].InputSchema, []string{"properties", "threads", "maxItems"}, float64(100)) + assertSchemaValue(t, tools[mcpcontract.ToolGetThreadFacets].InputSchema, []string{"properties", "facets", "maxItems"}, float64(10)) + assertSchemaValue(t, tools[mcpcontract.ToolGetThreadFacets].InputSchema, []string{"properties", "facets", "items", "enum"}, facets.AllNames()) assertSchemaValue(t, tools[mcpcontract.ToolHydrateThreads].InputSchema, []string{"properties", "max_pages", "default"}, float64(3)) assertSchemaValue(t, tools[mcpcontract.ToolRankThreads].InputSchema, []string{"required"}, []any{"repositories"}) assertSchemaValue(t, tools[mcpcontract.ToolCreateWorkspace].InputSchema, []string{"required"}, []any{"investigation_id"}) @@ -332,7 +336,8 @@ func TestAgentToolSelectionProxy(t *testing.T) { {"Run a repeat stress validation group with concurrency and telemetry", mcpcontract.ToolRunValidation}, {"Stop a running durable job", mcpcontract.ToolCancelJob}, {"Poll several durable jobs together with structured progress", mcpcontract.ToolGetJob}, - {"Read stored facet coverage for several exact threads", mcpcontract.ToolGetCoverage}, + {"Read stored facet coverage for several exact threads", mcpcontract.ToolGetThreadFacets}, + {"Read repository and thread coverage across several targets", mcpcontract.ToolGetCoverage}, {"Compare contribution candidates with my authored pull requests for overlap", mcpcontract.ToolFindPortfolioOverlaps}, {"Link an authored pull request to a local opportunity", mcpcontract.ToolLinkPullRequest}, {"Rebuild and persist the repository dossier from the local corpus", mcpcontract.ToolBuildRepositoryDossier}, diff --git a/internal/mcpserver/resource_templates.go b/internal/mcpserver/resource_templates.go index 43f0354e..1608e2dd 100644 --- a/internal/mcpserver/resource_templates.go +++ b/internal/mcpserver/resource_templates.go @@ -19,6 +19,12 @@ func (s *Server) registerResourceTemplates() { {"gitcontribute://readiness/{opportunity_id}", "Readiness", "Local contribution readiness report"}, {"gitcontribute://lens/{name}", "Lens", "Saved lens definition"}, } + if _, ok := s.reader.(threadFacetResourceReader); ok { + templates = append(templates, resourceTemplateDefinition{ + template: "gitcontribute://thread/{owner}/{repo}/{kind}/{number}/facet/{facet}", + name: "Thread facet", description: "Persisted local thread facet payload", + }) + } if _, ok := s.reader.(FixPatternReader); ok { templates = append(templates, resourceTemplateDefinition{ template: "gitcontribute://fix-pattern-report/{job_id}", diff --git a/internal/mcpserver/resources.go b/internal/mcpserver/resources.go index 8b463b41..caa08abb 100644 --- a/internal/mcpserver/resources.go +++ b/internal/mcpserver/resources.go @@ -34,6 +34,10 @@ type pullRequestWorkflowResourceReader interface { CIJobLogResource(context.Context, string, string, int, int64) (map[string]any, error) } +type threadFacetResourceReader interface { + ThreadFacetResource(context.Context, string, string, string, int, string) (map[string]any, error) +} + func (s *Server) readResource(ctx context.Context, req *mcp.ReadResourceRequest) (*mcp.ReadResourceResult, error) { uri := req.Params.URI u, err := url.Parse(uri) @@ -76,6 +80,9 @@ func (s *Server) readResourceValue(ctx context.Context, req resourceRequest) (an case "dossier": return s.readDossierResource(ctx, req) case "thread": + if len(req.parts) == 6 && req.parts[4] == "facet" { + return s.readThreadFacetResource(ctx, req) + } return s.readTypedThreadResource(ctx, req) case "investigation": return s.readInvestigationResource(ctx, req) @@ -108,6 +115,18 @@ func (s *Server) readResourceValue(ctx context.Context, req resourceRequest) (an } } +func (s *Server) readThreadFacetResource(ctx context.Context, req resourceRequest) (map[string]any, error) { + reader, ok := s.reader.(threadFacetResourceReader) + if !ok || len(req.parts) != 6 || req.parts[4] != "facet" { + return nil, mcp.ResourceNotFoundError(req.uri) + } + number, valid := positivePathNumber(req.parts[3]) + if !valid || strings.TrimSpace(req.parts[2]) == "" || strings.TrimSpace(req.parts[5]) == "" { + return nil, mcp.ResourceNotFoundError(req.uri) + } + return reader.ThreadFacetResource(ctx, req.parts[0], req.parts[1], req.parts[2], number, req.parts[5]) +} + func (s *Server) readPullRequestFeedbackResource(ctx context.Context, req resourceRequest) (map[string]any, error) { reader, ok := s.reader.(pullRequestWorkflowResourceReader) number, valid := pullRequestResourceNumber(req.parts) diff --git a/internal/mcpserver/scalable.go b/internal/mcpserver/scalable.go index adcec8ac..8148954e 100644 --- a/internal/mcpserver/scalable.go +++ b/internal/mcpserver/scalable.go @@ -19,8 +19,9 @@ const serverInstructions = "Use advertised GitContribute tools for durable, sour "The durable workflow is concern to investigation to hypothesis to opportunity to workspace to draft; use only advertised stages. " + "Use workflow.prepare_issue_set when exact issue numbers already define the contribution scope. " + "When an operation returns a job, poll advertised job tools in batches. " + + "Use corpus.get_thread_facets for bounded stored facet coverage and resources/read for larger facet payloads; repository, thread, and facet gaps provide the exact ordered synchronization route. " + "To inspect a returned resource, ask the host to perform MCP resources/read with this server and the exact URI; in Codex, call read_mcp_resource. Treat resource URIs as opaque identifiers and never shorten, pluralize, or reconstruct them. " + - "Missing or truncated coverage is unknown, not negative evidence; retry only retryable batch items. " + + "Missing or truncated coverage is unknown, not negative evidence; use each recovery plan's ordered typed calls and retry only retryable batch items. " + "Only advertised tools are available. GitContribute never mutates GitHub." // RepositoryRef identifies one GitHub repository without implying that it has @@ -114,11 +115,16 @@ const serverInstructions = "Use advertised GitContribute tools for durable, sour func (s *Server) registerScalable() { readOnly := readOnlyAnnotations() addCatalogTool(s, catalogTool[mcpcontract.GetRepositoriesInput, mcpcontract.GetRepositoriesOutput]{name: mcpcontract.ToolGetRepositories, title: "Get stored repositories in one batch", description: "Read metadata, coverage, and dossier availability for up to 100 stored repositories. Use for comparison before reading dossier resources. Missing metadata includes a sync action. Offline.", annotations: readOnly, supportedBy: supports[ScalableReader], input: inputSchema[mcpcontract.GetRepositoriesInput](func(sc *schemaBuilder) { setArrayBounds(sc, "repositories", 1, 100) }), output: outputSchema[mcpcontract.GetRepositoriesOutput]("Ordered repository batch with item-level status and dossier availability."), handler: s.getRepositories}) - addCatalogTool(s, catalogTool[mcpcontract.GetThreadsInput, mcpcontract.GetThreadsOutput]{name: mcpcontract.ToolGetThreads, title: "Get stored threads in one batch", description: "Read up to 100 exact stored issues or pull requests in input order. Choose compact for triage and full only for finalists; this tool is offline.", annotations: readOnly, supportedBy: supports[ScalableReader], input: inputSchema[mcpcontract.GetThreadsInput](func(sc *schemaBuilder) { + addCatalogTool(s, catalogTool[mcpcontract.GetThreadsInput, mcpcontract.GetThreadsOutput]{name: mcpcontract.ToolGetThreads, title: "Get stored threads in one batch", description: "Read exact stored issue or pull-request headers and optional complete bodies for up to 100 inputs. Choose compact for triage and full for finalist body reads; this tool is offline.", annotations: readOnly, supportedBy: supports[ScalableReader], input: inputSchema[mcpcontract.GetThreadsInput](func(sc *schemaBuilder) { setArrayBounds(sc, "threads", 1, 100) setEnum(sc, "view", "compact", "full") setDefault(sc, "view", "compact") }), output: outputSchema[mcpcontract.GetThreadsOutput]("Ordered stored-thread batch with item-level status."), handler: s.getThreads}) + addCatalogTool(s, catalogTool[mcpcontract.GetThreadFacetsInput, mcpcontract.GetThreadFacetsOutput]{name: mcpcontract.ToolGetThreadFacets, title: "Get facet coverage in one batch", description: "Read bounded offline coverage metadata and exact resource URIs for up to 100 thread facet selections. Use resources/read for persisted facet observations; this never contacts GitHub.", annotations: readOnly, supportedBy: supports[ThreadFacetReader], input: inputSchema[mcpcontract.GetThreadFacetsInput](func(sc *schemaBuilder) { + setArrayBounds(sc, "threads", 1, 100) + setArrayBounds(sc, "facets", 1, 10) + setArrayEnum(sc, "facets", facets.AllNames()...) + }), output: outputSchema[mcpcontract.GetThreadFacetsOutput]("Ordered stored-facet metadata with canonical resource links."), handler: s.getThreadFacets}) addCatalogTool(s, catalogTool[mcpcontract.RankOpportunitiesInput, mcpcontract.RankOpportunitiesOutput]{name: mcpcontract.ToolRankThreads, title: "Rank stored threads for contribution", description: "Rank open issues across 1-50 required stored repositories. This bounded offline result reports truncation and never persists opportunities.", annotations: readOnly, supportedBy: supports[ScalableReader], input: inputSchema[mcpcontract.RankOpportunitiesInput](func(sc *schemaBuilder) { setArrayBounds(sc, "repositories", 1, 50) setRange(sc, "limit", 1, 100) @@ -181,7 +187,7 @@ func (s *Server) registerScalable() { setDefault(sc, "max_requests", 1000) configureSyncThreadModes(sc) }), output: outputSchema[mcpcontract.JobReference]("Reference to a bounded thread-header synchronization job."), handler: s.syncThreads}) - addCatalogTool(s, catalogTool[mcpcontract.HydrateThreadsInput, mcpcontract.JobReference]{name: mcpcontract.ToolHydrateThreads, title: "Synchronize selected GitHub thread details", description: "Fetch selected GitHub child data for up to 100 known issues or pull requests. Use after ranking finalists. Do not use to inspect existing corpus coverage; corpus.get_coverage is the offline coverage read.", annotations: networkReadAnnotations(), supportedBy: supports[GitHubOperator], input: inputSchema[mcpcontract.HydrateThreadsInput](func(sc *schemaBuilder) { + addCatalogTool(s, catalogTool[mcpcontract.HydrateThreadsInput, mcpcontract.JobReference]{name: mcpcontract.ToolHydrateThreads, title: "Synchronize selected GitHub thread details", description: "Fetch selected GitHub child data for up to 100 known issues or pull requests. Use after ranking finalists. Do not use to inspect existing corpus coverage; corpus.get_thread_facets is the offline facet read and resources/read is the large-payload surface.", annotations: networkReadAnnotations(), supportedBy: supports[GitHubOperator], input: inputSchema[mcpcontract.HydrateThreadsInput](func(sc *schemaBuilder) { setArrayBounds(sc, "threads", 1, 100) setArrayBounds(sc, "facets", 1, 5) setArrayEnum(sc, "facets", facets.SelectableNames()...) @@ -328,6 +334,19 @@ func (s *Server) getThreads(ctx context.Context, _ *mcp.CallToolRequest, in mcpc out, err := r.GetThreads(ctx, in) return nil, out, err } +func (s *Server) getThreadFacets(ctx context.Context, _ *mcp.CallToolRequest, in mcpcontract.GetThreadFacetsInput) (*mcp.CallToolResult, mcpcontract.GetThreadFacetsOutput, error) { + for _, thread := range in.Threads { + if err := validateThreadRef(thread, false); err != nil { + return nil, mcpcontract.GetThreadFacetsOutput{}, err + } + } + r, ok := s.reader.(ThreadFacetReader) + if !ok { + return nil, mcpcontract.GetThreadFacetsOutput{}, errors.New("thread facet reads are not available") + } + out, err := r.GetThreadFacets(ctx, in) + return nil, out, err +} func (s *Server) rankOpportunities(ctx context.Context, _ *mcp.CallToolRequest, in mcpcontract.RankOpportunitiesInput) (*mcp.CallToolResult, mcpcontract.RankOpportunitiesOutput, error) { if in.Limit == 0 { in.Limit = 20 @@ -380,7 +399,11 @@ func (s *Server) getJobs(ctx context.Context, _ *mcp.CallToolRequest, in mcpcont normalizeJobExecution(&job) if in.ResponseFormat == "concise" { if job.Status == "succeeded" || job.Status == "failed" || job.Status == "cancelled" { - item.NextAction = "Call jobs.get with response_format=detailed to read typed artifact and follow-up references." + item.Recovery = &mcpcontract.RecoveryPlan{ + Version: mcpcontract.RecoveryPlanVersion, Reason: "blocked", + Message: "Read the detailed typed artifact and follow-up references.", + Then: []mcpcontract.ToolCall{{Tool: mcpcontract.ToolGetJob, Arguments: &mcpcontract.ToolCallArguments{IDs: []string{id}, ResponseFormat: "detailed"}}}, + } } } item.Value = &job @@ -453,7 +476,7 @@ func (s *Server) searchGitHubRepositories(ctx context.Context, _ *mcp.CallToolRe } out, err := op.SearchGitHubRepositories(ctx, in) if s.readOnly { - out.SuggestedActions = nil + out.RecoveryPlans = nil } return nil, out, err } diff --git a/internal/mcpserver/server.go b/internal/mcpserver/server.go index 8febf032..3b795617 100644 --- a/internal/mcpserver/server.go +++ b/internal/mcpserver/server.go @@ -30,6 +30,11 @@ type ScalableReader interface { GetJobs(context.Context, mcpcontract.GetJobsInput) (mcpcontract.GetJobsOutput, error) } +// ThreadFacetReader exposes the bounded offline facet metadata surface. +type ThreadFacetReader interface { + GetThreadFacets(context.Context, mcpcontract.GetThreadFacetsInput) (mcpcontract.GetThreadFacetsOutput, error) +} + // IssueSetReader prepares bounded contribution evidence from exact stored // issues without requiring or creating durable workflow state. type IssueSetReader interface { diff --git a/internal/mcpserver/server_contract_test.go b/internal/mcpserver/server_contract_test.go index b961970f..ba4773bd 100644 --- a/internal/mcpserver/server_contract_test.go +++ b/internal/mcpserver/server_contract_test.go @@ -24,6 +24,7 @@ func TestServerInstructionsContainRoutingPhrases(t *testing.T) { "explicit network reads", "concern to investigation to hypothesis to opportunity to workspace to draft", "poll advertised job tools in batches", + "recovery plan's ordered typed calls", "perform MCP resources/read", "in Codex, call read_mcp_resource", "exact URI", diff --git a/internal/mcpserver/server_test.go b/internal/mcpserver/server_test.go index fc6a16a2..50f1a0d3 100644 --- a/internal/mcpserver/server_test.go +++ b/internal/mcpserver/server_test.go @@ -53,11 +53,20 @@ func (*fakeReader) CIJobLogResource(context.Context, string, string, int, int64) return map[string]any{"schema_version": "gitcontribute.ci-job-log.v1", "body": "failure"}, nil } +func (*fakeReader) GetThreadFacets(_ context.Context, in mcpcontract.GetThreadFacetsInput) (mcpcontract.GetThreadFacetsOutput, error) { + return mcpcontract.GetThreadFacetsOutput{Status: "complete"}, nil +} + +func (*fakeReader) ThreadFacetResource(context.Context, string, string, string, int, string) (map[string]any, error) { + return map[string]any{"schema_version": "gitcontribute.thread-facet.v1"}, nil +} + func TestPullRequestWorkflowResourcesAreReadable(t *testing.T) { server := &Server{reader: &fakeReader{}} tests := []struct { uri, version string }{ + {"gitcontribute://thread/acme/project/pull_request/7/facet/pr_details", "gitcontribute.thread-facet.v1"}, {"gitcontribute://pull-request-feedback/acme/project/7", "gitcontribute.pull-request-feedback.v1"}, {"gitcontribute://ci-failure-report/acme/project/7", "gitcontribute.ci-failure-report.v1"}, {"gitcontribute://ci-job-log/acme/project/7/31", "gitcontribute.ci-job-log.v1"}, @@ -772,6 +781,7 @@ func TestV1ParityToolsAndResources(t *testing.T) { templates[template.URITemplate] = true } for _, uriTemplate := range []string{ + "gitcontribute://thread/{owner}/{repo}/{kind}/{number}/facet/{facet}", "gitcontribute://concern/{id}", "gitcontribute://draft/{id}/{revision}", "gitcontribute://manifest/{id}",