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

290 lines
8.9 KiB
Go

package agentstate
import (
"context"
"encoding/json"
"errors"
"os"
"path/filepath"
"reflect"
"sync"
"testing"
"time"
"iop/packages/go/agentpolicy"
"iop/packages/go/agentprovider/cli/status"
"iop/packages/go/agentruntime"
"iop/packages/go/agenttask"
)
func TestStoreRoundTripAndStaleCAS(t *testing.T) {
path := filepath.Join(t.TempDir(), "state", "manager.json")
store, err := NewStore(path)
if err != nil {
t.Fatalf("NewStore: %v", err)
}
state, revision, err := store.Load(context.Background())
if err != nil {
t.Fatalf("initial Load: %v", err)
}
if revision != "0" || state.SchemaVersion != agenttask.StateSchemaVersion {
t.Fatalf("initial revision/schema = %q/%d", revision, state.SchemaVersion)
}
state.NextOrdinal = 7
committed, err := store.CompareAndSwap(context.Background(), revision, state)
if err != nil {
t.Fatalf("CompareAndSwap: %v", err)
}
if committed != "1" {
t.Fatalf("committed revision = %q, want 1", committed)
}
reopened, err := NewStore(path)
if err != nil {
t.Fatalf("NewStore reopen: %v", err)
}
got, gotRevision, err := reopened.Load(context.Background())
if err != nil {
t.Fatalf("reopened Load: %v", err)
}
if gotRevision != "1" || got.NextOrdinal != 7 {
t.Fatalf("reopened revision/state = %q/%d", gotRevision, got.NextOrdinal)
}
if _, err := reopened.CompareAndSwap(context.Background(), "0", got); !errors.Is(err, agenttask.ErrRevisionConflict) {
t.Fatalf("stale CompareAndSwap error = %v", err)
}
}
func TestStoreRejectsCorruptionWithoutOverwrite(t *testing.T) {
path := filepath.Join(t.TempDir(), "manager.json")
store, err := NewStore(path)
if err != nil {
t.Fatalf("NewStore: %v", err)
}
state, revision, err := store.Load(context.Background())
if err != nil {
t.Fatalf("Load: %v", err)
}
if _, err := store.CompareAndSwap(context.Background(), revision, state); err != nil {
t.Fatalf("CompareAndSwap: %v", err)
}
payload, err := os.ReadFile(path)
if err != nil {
t.Fatalf("ReadFile: %v", err)
}
var envelope map[string]any
if err := json.Unmarshal(payload, &envelope); err != nil {
t.Fatalf("Unmarshal: %v", err)
}
envelope["checksum"] = "tampered"
corrupt, err := json.Marshal(envelope)
if err != nil {
t.Fatalf("Marshal: %v", err)
}
if err := os.WriteFile(path, corrupt, 0o600); err != nil {
t.Fatalf("WriteFile: %v", err)
}
before, err := os.ReadFile(path)
if err != nil {
t.Fatalf("ReadFile before rejected CAS: %v", err)
}
if _, _, err := store.Load(context.Background()); !errors.Is(err, ErrCorruptState) {
t.Fatalf("corrupt Load error = %v", err)
}
if _, err := store.CompareAndSwap(context.Background(), "1", state); !errors.Is(err, ErrCorruptState) {
t.Fatalf("corrupt CompareAndSwap error = %v", err)
}
after, err := os.ReadFile(path)
if err != nil {
t.Fatalf("ReadFile after rejected CAS: %v", err)
}
if string(after) != string(before) {
t.Fatal("rejected CAS overwrote corrupt checkpoint evidence")
}
}
func TestStoreConcurrentCASIsSerialized(t *testing.T) {
path := filepath.Join(t.TempDir(), "manager.json")
const writers = 12
var wait sync.WaitGroup
errs := make(chan error, writers)
for range writers {
wait.Add(1)
go func() {
defer wait.Done()
store, err := NewStore(path)
if err != nil {
errs <- err
return
}
for {
state, revision, err := store.Load(context.Background())
if err != nil {
errs <- err
return
}
state.NextOrdinal++
if _, err := store.CompareAndSwap(context.Background(), revision, state); errors.Is(err, agenttask.ErrRevisionConflict) {
continue
} else if err != nil {
errs <- err
}
return
}
}()
}
wait.Wait()
close(errs)
for err := range errs {
t.Errorf("concurrent writer: %v", err)
}
store, _ := NewStore(path)
state, revision, err := store.Load(context.Background())
if err != nil {
t.Fatalf("final Load: %v", err)
}
if state.NextOrdinal != writers || revision != "12" {
t.Fatalf("final ordinal/revision = %d/%s, want %d/12", state.NextOrdinal, revision, writers)
}
}
func TestStorePersistsSealedQuotaObservation(t *testing.T) {
path := filepath.Join(t.TempDir(), "manager.json")
store, err := NewStore(path)
if err != nil {
t.Fatalf("NewStore: %v", err)
}
now := time.Date(2026, 7, 29, 1, 2, 3, 0, time.UTC)
attempt := agenttask.AttemptID("attempt-v1/4:work/1:1")
locator := agenttask.LocatorRecord{
Kind: agenttask.LocatorProcess, Opaque: "pid:42:start:9001", Revision: "process-r1",
ProjectID: "project", WorkspaceID: "workspace",
WorkUnitID: "work", AttemptID: attempt,
}
sealedObservation := agentpolicy.NormalizeAttemptObservation(
status.NormalizeQuotaSnapshot(
"provider",
"profile",
[]string{"overall"},
now,
&status.UsageStatus{DailyLimit: "50%"},
nil,
),
&agentruntime.Failure{Code: agentruntime.FailureCodeUnavailable, Retryable: true},
now,
time.Minute,
)
state := agenttask.ManagerState{
SchemaVersion: agenttask.StateSchemaVersion,
DeviceLease: &agenttask.LeaseRecord{
OwnerID: "daemon", Token: "device-token", ExpiresAt: now.Add(time.Minute),
},
WorkspaceLeases: map[agenttask.WorkspaceID]agenttask.LeaseRecord{
"workspace": {
OwnerID: "daemon", Token: "workspace-token", ExpiresAt: now.Add(time.Minute),
},
},
Projects: map[agenttask.ProjectID]agenttask.ProjectRecord{
"project": {
ProjectID: "project", WorkspaceID: "workspace",
Status: agenttask.ProjectStatusRunning,
Works: map[agenttask.WorkUnitID]agenttask.WorkRecord{
"work": {
Unit: agenttask.WorkUnit{ID: "work", MilestoneID: "milestone"},
State: agenttask.WorkStateDispatching, Attempt: 1, AttemptID: attempt,
Locators: map[agenttask.LocatorKind]agenttask.LocatorRecord{
agenttask.LocatorProcess: locator,
},
FailureBudgets: map[agenttask.FailureStage]agenttask.FailureBudgetRecord{
agenttask.FailureStageDispatch: {
Stage: agenttask.FailureStageDispatch, Consecutive: 2, Limit: 10,
LastCode: agenttask.BlockerInvocationFailed,
AttemptID: attempt, UpdatedAt: now,
},
},
AttemptObservations: []agenttask.AttemptObservationRecord{{
AttemptID: attempt,
Target: agenttask.ExecutionTarget{
ProviderID: "provider", ModelID: "model", ProfileID: "profile",
ProfileRevision: "profile-r1", ConfigRevision: "config-r1", Capacity: 1,
},
Observation: sealedObservation,
}},
},
},
},
},
}
_, revision, err := store.Load(context.Background())
if err != nil {
t.Fatalf("Load: %v", err)
}
if _, err := store.CompareAndSwap(context.Background(), revision, state); err != nil {
t.Fatalf("CompareAndSwap: %v", err)
}
got, _, err := store.Load(context.Background())
if err != nil {
t.Fatalf("Load committed state: %v", err)
}
if !reflect.DeepEqual(got, state) {
t.Fatalf("recovery state changed across disk round trip\ngot: %#v\nwant: %#v", got, state)
}
payload, err := os.ReadFile(path)
if err != nil {
t.Fatalf("ReadFile: %v", err)
}
var envelope diskEnvelope
if err := decodeOne(payload, &envelope); err != nil {
t.Fatalf("decode envelope: %v", err)
}
var durableState map[string]any
if err := json.Unmarshal(envelope.State, &durableState); err != nil {
t.Fatalf("decode state object: %v", err)
}
quota := durableState["Projects"].(map[string]any)["project"].(map[string]any)["Works"].(map[string]any)["work"].(map[string]any)["AttemptObservations"].([]any)[0].(map[string]any)["Observation"].(map[string]any)["Quota"].(map[string]any)
quota["state"] = "exhausted"
envelope.State, err = json.Marshal(durableState)
if err != nil {
t.Fatalf("encode tampered state: %v", err)
}
envelope.Checksum, err = stateChecksum(envelope.SchemaVersion, envelope.Revision, envelope.State)
if err != nil {
t.Fatalf("stateChecksum: %v", err)
}
payload, err = json.Marshal(envelope)
if err != nil {
t.Fatalf("encode tampered envelope: %v", err)
}
if err := os.WriteFile(path, payload, 0o600); err != nil {
t.Fatalf("WriteFile tampered state: %v", err)
}
if _, _, err := store.Load(context.Background()); !errors.Is(err, ErrCorruptState) {
t.Fatalf("Load checksum-valid seal drift error = %v, want corrupt state", err)
}
}
func TestStoreRejectsUnsupportedEnvelopeSchema(t *testing.T) {
path := filepath.Join(t.TempDir(), "manager.json")
state, err := json.Marshal(agenttask.ManagerState{SchemaVersion: agenttask.StateSchemaVersion})
if err != nil {
t.Fatalf("Marshal state: %v", err)
}
checksum, err := stateChecksum(99, 1, state)
if err != nil {
t.Fatalf("stateChecksum: %v", err)
}
payload, err := json.Marshal(diskEnvelope{
SchemaVersion: 99, Revision: 1, Checksum: checksum, State: state,
})
if err != nil {
t.Fatalf("Marshal envelope: %v", err)
}
if err := os.WriteFile(path, payload, 0o600); err != nil {
t.Fatalf("WriteFile: %v", err)
}
store, _ := NewStore(path)
if _, _, err := store.Load(context.Background()); !errors.Is(err, ErrUnsupportedSchema) {
t.Fatalf("Load error = %v, want unsupported schema", err)
}
}