package taskloop import ( "bytes" "context" "errors" "fmt" "io" "os" "path/filepath" "reflect" "strings" "iop/packages/go/agentconfig" "iop/packages/go/agenttask" ) const maxWorkflowArtifactBytes = 2 << 20 // Evidence applies the same bounded artifact gate to every provider. Pi may // receive a repair intent, but repair still requires an explicitly configured // native-context executor; this adapter never edits review evidence itself. type Evidence struct { snapshot agentconfig.RuntimeSnapshot roots ArtifactRootResolver repair EvidenceRepairer } type EvidenceRepairer interface { RepairEvidence(context.Context, agenttask.WorkflowEvidenceRepairRequest) error } // PiEvidenceRepairer resumes only the exact native session locator persisted // for the retained work attempt. It never edits evidence directly; Manager // performs a fresh Observe after this confined continuation returns. type PiEvidenceRepairer struct { catalog agentconfig.Catalog backend retainedIsolationBackend } func NewPiEvidenceRepairer( catalogConfig agentconfig.Catalog, backend retainedIsolationBackend, ) (*PiEvidenceRepairer, error) { if backend == nil { return nil, errors.New("taskloop: Pi evidence repair requires a retained isolation backend") } normalized, err := agentconfig.Normalize(catalogConfig) if err != nil { return nil, err } return &PiEvidenceRepairer{catalog: normalized, backend: backend}, nil } var _ agenttask.WorkflowEvidence = (*Evidence)(nil) func NewEvidence( snapshot agentconfig.RuntimeSnapshot, roots ArtifactRootResolver, repair EvidenceRepairer, ) *Evidence { return &Evidence{snapshot: snapshot, roots: roots, repair: repair} } func (evidence *Evidence) Observe( ctx context.Context, request agenttask.WorkflowEvidenceRequest, ) (agenttask.ArtifactEvidence, error) { if err := ctx.Err(); err != nil { return agenttask.ArtifactEvidence{}, err } registration, ok := evidence.snapshot.Project(string(request.Project.ProjectID)) if !ok { return agenttask.ArtifactEvidence{}, errors.New("taskloop: evidence project is not registered") } canonicalRoot, err := canonicalDirectory(registration.Workspace) if err != nil { return agenttask.ArtifactEvidence{}, err } if WorkspaceIdentity(canonicalRoot) != request.Project.WorkspaceID { return agenttask.ArtifactEvidence{}, errors.New("taskloop: evidence workspace identity mismatch") } if evidence.roots == nil { return agenttask.ArtifactEvidence{}, errors.New("taskloop: retained artifact root resolver is unavailable") } artifactRoot, err := evidence.roots.ArtifactRoot(request.Work) if err != nil { return agenttask.ArtifactEvidence{}, err } reviewPath := request.Work.Unit.Metadata["review_path"] if reviewPath == "" { return agenttask.ArtifactEvidence{}, errors.New("taskloop: review artifact locator is missing") } inspected, err := inspectReviewArtifact(artifactRoot, reviewPath) if err != nil { if errors.Is(err, os.ErrNotExist) { return agenttask.ArtifactEvidence{Active: false}, nil } return agenttask.ArtifactEvidence{}, err } identity := agenttask.ArtifactIdentity{ ProjectID: request.Project.ProjectID, WorkspaceID: request.Project.WorkspaceID, WorkUnitID: request.Work.Unit.ID, AttemptID: request.Work.AttemptID, ArtifactID: request.Submission.ArtifactID, } result := agenttask.ArtifactEvidence{ Active: true, Identity: identity, } if !inspected.placeholder { result.Completeness = agenttask.ArtifactComplete return result, nil } result.Completeness = agenttask.ArtifactPlaceholder if request.Work.Target != nil && strings.EqualFold(request.Work.Target.ProviderID, "pi") { if locator, ok := request.Work.Locators[agenttask.LocatorSession]; ok { result.RepairIntent = &agenttask.EvidenceRepairIntent{ Identity: identity, NativeLocator: locator, DispatchOrdinal: request.Work.DispatchOrdinal, } } } return result, nil } func (evidence *Evidence) Repair( ctx context.Context, request agenttask.WorkflowEvidenceRepairRequest, ) error { if evidence.repair == nil { return errors.New("taskloop: Pi evidence repair executor is not configured") } return evidence.repair.RepairEvidence(ctx, request) } func (repairer *PiEvidenceRepairer) RepairEvidence( ctx context.Context, request agenttask.WorkflowEvidenceRepairRequest, ) error { if err := ctx.Err(); err != nil { return err } if request.Work.Target == nil || !strings.EqualFold(request.Work.Target.ProviderID, "pi") { return errors.New("taskloop: evidence repair is restricted to Pi") } persisted, ok := request.Work.Locators[agenttask.LocatorSession] if !ok || persisted.Kind != agenttask.LocatorSession || !reflect.DeepEqual(persisted, request.Intent.NativeLocator) { return errors.New("taskloop: Pi repair native session locator mismatch") } expectedIdentity := agenttask.ArtifactIdentity{ ProjectID: request.Project.ProjectID, WorkspaceID: request.Project.WorkspaceID, WorkUnitID: request.Work.Unit.ID, AttemptID: request.Work.AttemptID, ArtifactID: request.Submission.ArtifactID, } if !reflect.DeepEqual(request.Intent.Identity, expectedIdentity) || request.Intent.DispatchOrdinal != request.Work.DispatchOrdinal || persisted.ProjectID != request.Project.ProjectID || persisted.WorkspaceID != request.Project.WorkspaceID || persisted.WorkUnitID != request.Work.Unit.ID || persisted.AttemptID != request.Work.AttemptID || persisted.Revision != digestStrings( "native-session-locator", persisted.Opaque, ) { return errors.New("taskloop: Pi repair intent identity mismatch") } var reference nativeSessionReference if err := decodeStrictJSON([]byte(persisted.Opaque), &reference); err != nil { return fmt.Errorf("taskloop: decode Pi native session locator: %w", err) } prepared, descriptor, err := rehydrateRetainedIsolation( ctx, repairer.backend, repairer.catalog, request.Project, request.Work, ) if err != nil { return err } resolved, ok := repairer.catalog.ResolveProfile(request.Work.Target.ProfileID) if !ok { return errors.New("taskloop: Pi repair profile is not declared") } binding := prepared.Confinement.Binding() if err := validateNativeSessionReference( reference, resolved, *request.Work.Target, binding, ); err != nil { return err } nativeSessionFile, err := resolveNativeSessionFile(reference) if err != nil { return err } planRelative := request.Work.Unit.Metadata["plan_path"] reviewRelative := request.Work.Unit.Metadata["review_path"] if planRelative == "" || reviewRelative == "" { return errors.New("taskloop: Pi repair artifact locators are incomplete") } prompt := fmt.Sprintf( "Continue this exact native session for %s. Complete only the "+ "implementation-owned evidence fields in %s, keep artifact content "+ "in English, and stop ready for official review.", planRelative, reviewRelative, ) command, err := prepareResumeCatalogCommand( resolved, *request.Work.Target, binding, descriptor.WorkingDir, prompt, nativeSessionFile, reference, ) if err != nil { return err } return invokeConfined(ctx, prepared.Confinement, command) } type reviewArtifact struct { path string placeholder bool } func inspectReviewArtifact(root, relative string) (reviewArtifact, error) { canonicalRoot, err := canonicalDirectory(root) if err != nil { return reviewArtifact{}, err } if filepath.IsAbs(relative) || filepath.Clean(relative) == "." || strings.HasPrefix(filepath.Clean(relative), ".."+string(filepath.Separator)) { return reviewArtifact{}, errors.New("taskloop: review artifact must be a contained relative path") } path := filepath.Join(canonicalRoot, filepath.FromSlash(relative)) canonicalParent, err := filepath.EvalSymlinks(filepath.Dir(path)) if err != nil { return reviewArtifact{}, err } relativeParent, err := filepath.Rel(canonicalRoot, canonicalParent) if err != nil || relativeParent == ".." || strings.HasPrefix(relativeParent, ".."+string(filepath.Separator)) { return reviewArtifact{}, errors.New("taskloop: review artifact escapes the registered workspace") } file, err := os.Open(path) if err != nil { return reviewArtifact{}, err } defer file.Close() info, err := file.Stat() if err != nil { return reviewArtifact{}, err } if !info.Mode().IsRegular() || info.Size() <= 0 || info.Size() > maxWorkflowArtifactBytes { return reviewArtifact{}, errors.New("taskloop: review artifact is not one bounded regular file") } content := make([]byte, info.Size()) if _, err := io.ReadFull(file, content); err != nil { return reviewArtifact{}, err } return reviewArtifact{ path: path, placeholder: hasImplementationPlaceholder(content), }, nil } func hasImplementationPlaceholder(content []byte) bool { for _, line := range bytes.Split(content, []byte{'\n'}) { text := strings.TrimSpace(string(line)) if strings.HasPrefix(text, "_Paste actual ") || text == "_Record actual deviations or `None`._" || text == "_Record actual decisions._" || strings.Contains(text, "[ ] Fill implementation-owned sections") { return true } } return false } func artifactPath(root, relative string) (string, error) { artifact, err := inspectReviewArtifact(root, relative) if err != nil { return "", fmt.Errorf("taskloop: inspect review artifact: %w", err) } return artifact.path, nil }