독립 호스트에서 안전한 작업 실행과 복구를 제공하기 위해 런타임 설정, 정책, 상태 저장소, 워크스페이스 격리 및 AgentTask 오케스트레이션을 확장한다.
290 lines
8.9 KiB
Go
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)
|
|
}
|
|
}
|