iop/packages/go/agenttask/workflow_evidence.go
toki 3c48879c44 feat(agent-runtime): 독립 Agent CLI 런타임 기반을 구현한다
독립 호스트에서 안전한 작업 실행과 복구를 제공하기 위해 런타임 설정, 정책, 상태 저장소, 워크스페이스 격리 및 AgentTask 오케스트레이션을 확장한다.
2026-07-29 13:12:21 +09:00

146 lines
5.1 KiB
Go

package agenttask
import (
"context"
"fmt"
)
// WorkflowEvidenceMatch is the provider-neutral result of comparing a
// project adapter observation to the manager's durable submission identity.
type WorkflowEvidenceMatch struct {
Status WorkflowEvidenceStatus
Identity ArtifactIdentity
RepairIntent *EvidenceRepairIntent
}
func artifactIdentity(project ProjectRecord, work WorkRecord, submission Submission) ArtifactIdentity {
return ArtifactIdentity{
ProjectID: project.ProjectID, WorkspaceID: project.WorkspaceID,
WorkUnitID: work.Unit.ID, AttemptID: work.AttemptID, ArtifactID: submission.ArtifactID,
}
}
// MatchWorkflowEvidence keeps project-specific artifact parsing outside the
// common package while enforcing exact active-pair identity and completeness.
func MatchWorkflowEvidence(expected ArtifactIdentity, observed ArtifactEvidence) WorkflowEvidenceMatch {
match := WorkflowEvidenceMatch{Identity: observed.Identity}
if !observed.Active {
match.Status = WorkflowEvidenceNotActive
return match
}
if observed.Identity != expected {
match.Status = WorkflowEvidenceIdentityMismatch
return match
}
switch observed.Completeness {
case ArtifactComplete:
match.Status = WorkflowEvidenceMatched
case ArtifactPlaceholder:
match.Status = WorkflowEvidencePlaceholder
match.RepairIntent = observed.RepairIntent
default:
match.Status = WorkflowEvidenceInvalid
}
return match
}
func (m *Manager) gateSubmissionEvidence(
ctx context.Context,
project ProjectRecord,
work WorkRecord,
submission Submission,
) *Blocker {
request := WorkflowEvidenceRequest{Project: project, Work: work, Submission: submission}
observed, err := m.evidence.Observe(ctx, request)
if err != nil {
return &Blocker{
Code: BlockerEvidenceUnavailable, Message: fmt.Sprintf("workflow evidence observation failed: %v", err),
Retryable: true,
}
}
match := MatchWorkflowEvidence(artifactIdentity(project, work, submission), observed)
if match.Status == WorkflowEvidenceMatched {
return nil
}
if match.Status != WorkflowEvidencePlaceholder {
return blockerForWorkflowEvidence(match.Status, "")
}
if work.Target == nil || work.Target.ProviderID != "pi" {
return &Blocker{
Code: BlockerEvidenceRepairDenied,
Message: "only Pi may repair a placeholder workflow artifact",
}
}
if match.RepairIntent == nil {
return &Blocker{
Code: BlockerEvidenceRepairDenied,
Message: "Pi placeholder workflow artifact has no same-context repair intent",
}
}
if err := validateEvidenceRepairIntent(project, work, submission, *match.RepairIntent); err != nil {
return &Blocker{Code: BlockerEvidenceRepairDenied, Message: err.Error()}
}
if err := m.evidence.Repair(ctx, WorkflowEvidenceRepairRequest{
Project: project, Work: work, Submission: submission, Intent: *match.RepairIntent,
}); err != nil {
return &Blocker{
Code: BlockerEvidenceRepairFailed, Message: fmt.Sprintf("Pi workflow evidence repair failed: %v", err),
Retryable: true,
}
}
// A repair never authorizes review by itself. Observe again and require a
// fresh exact match of the same active pair before accepting submission.
rematchedEvidence, err := m.evidence.Observe(ctx, request)
if err != nil {
return &Blocker{
Code: BlockerEvidenceUnavailable, Message: fmt.Sprintf("workflow evidence rematch failed: %v", err),
Retryable: true,
}
}
rematch := MatchWorkflowEvidence(artifactIdentity(project, work, submission), rematchedEvidence)
if rematch.Status == WorkflowEvidenceMatched {
return nil
}
return blockerForWorkflowEvidence(rematch.Status, "Pi repair did not produce a matched workflow artifact")
}
func blockerForWorkflowEvidence(status WorkflowEvidenceStatus, prefix string) *Blocker {
detail := string(status)
if prefix != "" {
detail = prefix + ": " + detail
}
switch status {
case WorkflowEvidenceIdentityMismatch:
return &Blocker{Code: BlockerArtifactMismatch, Message: detail}
case WorkflowEvidenceNotActive, WorkflowEvidencePlaceholder, WorkflowEvidenceInvalid:
return &Blocker{Code: BlockerSubmissionIncomplete, Message: detail}
default:
return &Blocker{Code: BlockerEvidenceUnavailable, Message: detail}
}
}
func validateEvidenceRepairIntent(
project ProjectRecord,
work WorkRecord,
submission Submission,
intent EvidenceRepairIntent,
) error {
if intent.Identity != artifactIdentity(project, work, submission) {
return fmt.Errorf("workflow evidence repair identity does not match the active submission")
}
if intent.DispatchOrdinal == 0 || intent.DispatchOrdinal != work.DispatchOrdinal {
return fmt.Errorf("workflow evidence repair ordinal does not match the active dispatch")
}
if intent.NativeLocator.Kind != LocatorSession {
return fmt.Errorf("workflow evidence repair requires a durable native session locator")
}
if err := validateLocator(project, work, intent.NativeLocator); err != nil {
return fmt.Errorf("workflow evidence repair locator is invalid: %w", err)
}
persisted, ok := work.Locators[LocatorSession]
if !ok || persisted != intent.NativeLocator {
return fmt.Errorf("workflow evidence repair locator is stale or was not durably checkpointed")
}
return nil
}