Node의 provider progress 기반 stall timeout, watchdog fencing과 bounded health probe evidence를 실행 경로에 반영한다. Edge-Node 계약과 구현 스펙, 테스트 및 Milestone 완료 evidence를 현재 상태와 맞춘다.
382 lines
16 KiB
Go
382 lines
16 KiB
Go
package node
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
"google.golang.org/protobuf/proto"
|
|
|
|
"iop/packages/go/credentiallease"
|
|
runtime "iop/packages/go/execution"
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
// TestRunWatchdogLifecycle covers confirmed fence with exact grace, unconfirmed
|
|
// ownership retention until provider return, caller cancel winning the timer
|
|
// race, and an already-expired deadline bypassing the watchdog.
|
|
func TestRunWatchdogLifecycle(t *testing.T) {
|
|
t.Run("confirmed fence and exact grace", testRunWatchdogConfirmed)
|
|
t.Run("unconfirmed retains ownership until provider return", testRunWatchdogUnconfirmed)
|
|
t.Run("caller cancel wins timer race", testRunWatchdogCancelPrecedence)
|
|
t.Run("hard deadline retains boundary", testRunWatchdogDeadlinePrecedence)
|
|
}
|
|
|
|
func testRunWatchdogConfirmed(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("run-confirmed")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnRunRequest(context.Background(), pipe.sess, &iop.RunRequest{RunId: "run-confirmed", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000})
|
|
}()
|
|
call := <-adapter.runCalls
|
|
stallTimer := clock.waitTimer(t, 0)
|
|
requireTimerDurations(t, stallTimer, time.Second)
|
|
stallTimer.fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
grace := clock.waitTimer(t, 1)
|
|
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
|
|
adapter.runReturn <- nil
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("run result = %v", err)
|
|
}
|
|
event := waitRunEvent(t, pipe.events)
|
|
if event.GetType() != string(runtime.EventTypeError) || event.GetMetadata()["attempt_fence"] != "confirmed" || event.GetMetadata()["idle_duration_ms"] != "1000" {
|
|
t.Fatalf("stall event = %+v", event)
|
|
}
|
|
if activeAdapterAttempts(n, adapter.Name()) != 0 || n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("confirmed provider return retained local ownership")
|
|
}
|
|
_ = call.sink.Emit(context.Background(), runtime.RuntimeEvent{RunID: "run-confirmed", Type: runtime.EventTypeDelta, Delta: "late"})
|
|
select {
|
|
case extra := <-pipe.events:
|
|
t.Fatalf("late or duplicate event = %+v", extra)
|
|
default:
|
|
}
|
|
}
|
|
|
|
func testRunWatchdogUnconfirmed(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("run-unconfirmed")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnRunRequest(context.Background(), pipe.sess, &iop.RunRequest{RunId: "run-unconfirmed", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 2000})
|
|
}()
|
|
call := <-adapter.runCalls
|
|
stallTimer := clock.waitTimer(t, 0)
|
|
_ = call.sink.Emit(context.Background(), runtime.RuntimeEvent{RunID: "run-unconfirmed", Type: runtime.EventTypeDelta, Delta: "progress"})
|
|
if progress := waitRunEvent(t, pipe.events); progress.GetType() != string(runtime.EventTypeDelta) {
|
|
t.Fatalf("progress event = %+v", progress)
|
|
}
|
|
requireTimerDurations(t, stallTimer, 2*time.Second, 2*time.Second)
|
|
stallTimer.fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
grace := clock.waitTimer(t, 1)
|
|
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
|
|
grace.fire()
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("run result = %v", err)
|
|
}
|
|
if event := waitRunEvent(t, pipe.events); event.GetMetadata()["attempt_fence"] != "unconfirmed" {
|
|
t.Fatalf("stall event = %+v", event)
|
|
}
|
|
if activeAdapterAttempts(n, adapter.Name()) != 1 || !n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("unconfirmed attempt released ownership before provider return")
|
|
}
|
|
_ = call.sink.Emit(context.Background(), runtime.RuntimeEvent{RunID: "run-unconfirmed", Type: runtime.EventTypeDelta, Delta: "late"})
|
|
adapter.runReturn <- nil
|
|
waitForOwnershipRelease(t, n, adapter.Name(), "provider return did not release retained ownership")
|
|
select {
|
|
case extra := <-pipe.events:
|
|
t.Fatalf("late or duplicate event = %+v", extra)
|
|
default:
|
|
}
|
|
}
|
|
|
|
func testRunWatchdogCancelPrecedence(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("run-cancel")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnRunRequest(ctx, pipe.sess, &iop.RunRequest{RunId: "run-cancel", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000})
|
|
}()
|
|
call := <-adapter.runCalls
|
|
timer := clock.waitTimer(t, 0)
|
|
cancel()
|
|
waitContextCanceled(t, call.ctx)
|
|
timer.fire()
|
|
adapter.runReturn <- runtime.ErrRunCancelled
|
|
if err := <-done; err != runtime.ErrRunCancelled {
|
|
t.Fatalf("cancel result = %v", err)
|
|
}
|
|
event := waitRunEvent(t, pipe.events)
|
|
if event.GetType() != string(runtime.EventTypeCancelled) || event.GetMetadata()["failure_code"] == string(runtime.FailureCodeResponseStalled) {
|
|
t.Fatalf("cancel event relabeled as stall: %+v", event)
|
|
}
|
|
if clock.count() != 1 {
|
|
t.Fatalf("cancel created close-grace timer: %d timers", clock.count())
|
|
}
|
|
}
|
|
|
|
func testRunWatchdogDeadlinePrecedence(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("run-deadline")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
ctx, cancel := context.WithDeadline(context.Background(), time.Now().Add(-time.Second))
|
|
defer cancel()
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnRunRequest(ctx, pipe.sess, &iop.RunRequest{RunId: "run-deadline", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000})
|
|
}()
|
|
call := <-adapter.runCalls
|
|
stallTimer := clock.waitTimer(t, 0)
|
|
waitContextCanceled(t, call.ctx)
|
|
adapter.runReturn <- context.DeadlineExceeded
|
|
if err := <-done; err != context.DeadlineExceeded {
|
|
t.Fatalf("deadline result = %v", err)
|
|
}
|
|
stallTimer.fire()
|
|
event := waitRunEvent(t, pipe.events)
|
|
if event.GetType() != string(runtime.EventTypeError) || event.GetError() != context.DeadlineExceeded.Error() || event.GetMetadata()["failure_code"] == string(runtime.FailureCodeResponseStalled) {
|
|
t.Fatalf("deadline event relabeled as stall: %+v", event)
|
|
}
|
|
if clock.count() != 1 {
|
|
t.Fatalf("deadline created close-grace timer: %d timers", clock.count())
|
|
}
|
|
}
|
|
|
|
// TestTunnelWatchdogLifecycle covers unconfirmed fence dropping late frames,
|
|
// confirmed fence, provider terminal stopping the clock, and credential
|
|
// ownership following provider return.
|
|
func TestTunnelWatchdogLifecycle(t *testing.T) {
|
|
t.Run("unconfirmed fence drops late frames", testTunnelWatchdogUnconfirmed)
|
|
t.Run("confirmed fence", testTunnelWatchdogConfirmed)
|
|
t.Run("provider terminal stops clock", testTunnelProviderTerminalStopsClock)
|
|
t.Run("credential ownership follows provider return", testTunnelCredentialOwnership)
|
|
}
|
|
|
|
func testTunnelWatchdogUnconfirmed(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("tunnel-unconfirmed")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnProviderTunnelRequest(context.Background(), pipe.sess, &iop.ProviderTunnelRequest{RunId: "tunnel-run", TunnelId: "tunnel", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1500})
|
|
}()
|
|
call := <-adapter.tunnelCalls
|
|
stallTimer := clock.waitTimer(t, 0)
|
|
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{RunID: "tunnel-run", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindBody, Body: []byte("progress")}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if frame := waitTunnelFrame(t, pipe.frames); frame.GetKind() != iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_BODY {
|
|
t.Fatalf("progress frame = %+v", frame)
|
|
}
|
|
requireTimerDurations(t, stallTimer, 1500*time.Millisecond, 1500*time.Millisecond)
|
|
stallTimer.fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
grace := clock.waitTimer(t, 1)
|
|
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
|
|
grace.fire()
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("tunnel result = %v", err)
|
|
}
|
|
terminal := waitTunnelFrame(t, pipe.frames)
|
|
if terminal.GetKind() != iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR || terminal.GetMetadata()["attempt_fence"] != "unconfirmed" {
|
|
t.Fatalf("stall terminal = %+v", terminal)
|
|
}
|
|
if activeAdapterAttempts(n, adapter.Name()) != 1 || !n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("unconfirmed tunnel released ownership before provider return")
|
|
}
|
|
_ = call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{RunID: "tunnel-run", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindUsage})
|
|
adapter.tunnelReturn <- nil
|
|
waitForOwnershipRelease(t, n, adapter.Name(), "tunnel provider return did not release ownership")
|
|
select {
|
|
case extra := <-pipe.frames:
|
|
t.Fatalf("late or duplicate frame = %+v", extra)
|
|
default:
|
|
}
|
|
}
|
|
|
|
func testTunnelWatchdogConfirmed(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("tunnel-confirmed")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnProviderTunnelRequest(context.Background(), pipe.sess, &iop.ProviderTunnelRequest{RunId: "tunnel-confirmed-run", TunnelId: "tunnel-confirmed", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000})
|
|
}()
|
|
call := <-adapter.tunnelCalls
|
|
clock.waitTimer(t, 0).fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
grace := clock.waitTimer(t, 1)
|
|
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
|
|
adapter.tunnelReturn <- nil
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("tunnel result = %v", err)
|
|
}
|
|
if terminal := waitTunnelFrame(t, pipe.frames); terminal.GetMetadata()["attempt_fence"] != "confirmed" {
|
|
t.Fatalf("stall terminal = %+v", terminal)
|
|
}
|
|
if activeAdapterAttempts(n, adapter.Name()) != 0 || n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("confirmed tunnel retained ownership")
|
|
}
|
|
}
|
|
|
|
func testTunnelProviderTerminalStopsClock(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("tunnel-terminal")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
pipe := newWatchdogPipe(t)
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.OnProviderTunnelRequest(context.Background(), pipe.sess, &iop.ProviderTunnelRequest{RunId: "terminal-run", TunnelId: "terminal-tunnel", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000})
|
|
}()
|
|
call := <-adapter.tunnelCalls
|
|
timer := clock.waitTimer(t, 0)
|
|
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{RunID: "terminal-run", TunnelID: "terminal-tunnel", Kind: runtime.ProviderTunnelFrameKindEnd, End: true}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
_ = waitTunnelFrame(t, pipe.frames)
|
|
_, stopped := timer.snapshot()
|
|
if !stopped {
|
|
t.Fatal("provider terminal did not stop tunnel watchdog")
|
|
}
|
|
adapter.tunnelReturn <- nil
|
|
if err := <-done; err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func testTunnelCredentialOwnership(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("tunnel-credential")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
ticket, err := n.admissionFor(adapter.Name(), runtime.Capabilities{MaxConcurrency: 1}).acquire()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
secret := []byte("provider-secret")
|
|
credential := &runtime.ProviderCredential{HeaderName: "Authorization", Scheme: "Bearer", Secret: secret}
|
|
material := &credentiallease.Material{HeaderName: credential.HeaderName, Scheme: credential.Scheme, Secret: secret}
|
|
tr := runtime.ProviderTunnelRequest{RunID: "credential-run", TunnelID: "credential-tunnel", Adapter: adapter.Name(), Target: "target", Credential: credential, ResponseStallTimeoutMS: 1000}
|
|
execCtx, cancel := context.WithCancel(context.Background())
|
|
h := &runHandle{runID: tr.RunID, adapter: tr.Adapter, target: tr.Target, cancel: cancel, done: make(chan struct{})}
|
|
n.runs.register(h)
|
|
sink := &tunnelSink{sess: noopSender{}, observer: newAttemptObserver(clock, time.Second)}
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.executeTunnelAttempt(execCtx, cancel, adapter, tr, sink, ticket, h, material, nil, nil)
|
|
}()
|
|
call := <-adapter.tunnelCalls
|
|
clock.waitTimer(t, 0).fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
clock.waitTimer(t, 1).fire()
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("tunnel result = %v", err)
|
|
}
|
|
if string(credential.Secret) != "provider-secret" || string(material.Secret) != "provider-secret" || activeAdapterAttempts(n, adapter.Name()) != 1 || !n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("unconfirmed tunnel did not retain credential and local ownership")
|
|
}
|
|
adapter.tunnelReturn <- nil
|
|
select {
|
|
case <-h.done:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("credential cleanup did not follow provider return")
|
|
}
|
|
if credential.Secret != nil || material.Secret != nil || activeAdapterAttempts(n, adapter.Name()) != 0 || n.runs.hasAnyActiveRuns() {
|
|
t.Fatal("provider return did not zero credentials and release local ownership")
|
|
}
|
|
}
|
|
|
|
type tunnelTerminalOwnership struct {
|
|
admissionReleased bool
|
|
runDeregistered bool
|
|
credentialsZeroed bool
|
|
handleClosed bool
|
|
}
|
|
|
|
type tunnelTerminalInspector struct {
|
|
ownership func() tunnelTerminalOwnership
|
|
seen chan tunnelTerminalOwnership
|
|
}
|
|
|
|
func (s *tunnelTerminalInspector) Send(message proto.Message) error {
|
|
frame, ok := message.(*iop.ProviderTunnelFrame)
|
|
if ok && frame.GetKind() == iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR {
|
|
s.seen <- s.ownership()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// TestTunnelConfirmedFenceClosesOwnershipBeforeTerminal proves the confirmed
|
|
// terminal is visible to the edge only after admission, run deregistration,
|
|
// credential zeroing, and handle closure have all completed.
|
|
func TestTunnelConfirmedFenceClosesOwnershipBeforeTerminal(t *testing.T) {
|
|
clock := newManualAttemptClock()
|
|
adapter := newControlledWatchdogAdapter("tunnel-confirmed-ownership")
|
|
n := newWatchdogNode(t, adapter, clock)
|
|
ticket, err := n.admissionFor(adapter.Name(), runtime.Capabilities{MaxConcurrency: 1}).acquire()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
credential := &runtime.ProviderCredential{HeaderName: "Authorization", Scheme: "Bearer", Secret: []byte("provider-secret")}
|
|
material := &credentiallease.Material{HeaderName: credential.HeaderName, Scheme: credential.Scheme, Secret: []byte("provider-secret")}
|
|
tr := runtime.ProviderTunnelRequest{RunID: "tunnel-confirmed-ownership", TunnelID: "tunnel", Adapter: adapter.Name(), Target: "target", Credential: credential, ResponseStallTimeoutMS: 1000}
|
|
execCtx, cancel := context.WithCancel(context.Background())
|
|
h := &runHandle{runID: tr.RunID, adapter: tr.Adapter, target: tr.Target, cancel: cancel, done: make(chan struct{})}
|
|
n.runs.register(h)
|
|
inspector := &tunnelTerminalInspector{seen: make(chan tunnelTerminalOwnership, 1)}
|
|
inspector.ownership = func() tunnelTerminalOwnership {
|
|
ownership := tunnelTerminalOwnership{
|
|
admissionReleased: activeAdapterAttempts(n, adapter.Name()) == 0,
|
|
runDeregistered: !n.runs.hasAnyActiveRuns(),
|
|
credentialsZeroed: credential.Secret == nil && material.Secret == nil,
|
|
}
|
|
select {
|
|
case <-h.done:
|
|
ownership.handleClosed = true
|
|
default:
|
|
}
|
|
return ownership
|
|
}
|
|
sink := &tunnelSink{sess: inspector, observer: newAttemptObserver(clock, time.Second)}
|
|
done := make(chan error, 1)
|
|
go func() {
|
|
done <- n.executeTunnelAttempt(execCtx, cancel, adapter, tr, sink, ticket, h, material, nil, nil)
|
|
}()
|
|
call := <-adapter.tunnelCalls
|
|
clock.waitTimer(t, 0).fire()
|
|
waitContextCanceled(t, call.ctx)
|
|
grace := clock.waitTimer(t, 1)
|
|
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
|
|
adapter.tunnelReturn <- nil
|
|
if err := <-done; err != errProviderResponseStalled {
|
|
t.Fatalf("tunnel result = %v", err)
|
|
}
|
|
ownership := <-inspector.seen
|
|
if !ownership.admissionReleased || !ownership.runDeregistered || !ownership.credentialsZeroed || !ownership.handleClosed {
|
|
t.Fatalf("confirmed terminal was visible before local ownership closed: %+v", ownership)
|
|
}
|
|
}
|
|
|
|
func waitForOwnershipRelease(t *testing.T, n *Node, adapter, failure string) {
|
|
t.Helper()
|
|
deadline := time.After(2 * time.Second)
|
|
for activeAdapterAttempts(n, adapter) != 0 || n.runs.hasAnyActiveRuns() {
|
|
select {
|
|
case <-deadline:
|
|
t.Fatal(failure)
|
|
default:
|
|
}
|
|
}
|
|
}
|