// Independent CloseDeadline adapter; all protocol behavior is pinned upstream.
package main
import (
 "encoding/json"
 "errors"
 "net"
 "os"
 "strings"
 "sync"
 "time"
 amqp "github.com/rabbitmq/amqp091-go"
)
type Request struct{URI,Scenario string;Milliseconds int}
func check(e error){if e!=nil{panic(e)}}
func kind(e error)string{if e==nil{return "ok"};if errors.Is(e,amqp.ErrClosed){return "closed"};if strings.Contains(strings.ToLower(e.Error()),"timeout")||strings.Contains(e.Error(),"deadline"){return "timeout"};return "error"}
func cancelled(ch <-chan struct{})bool{select{case <-ch:return true;default:return false}}
func main(){
 var r Request;check(json.NewDecoder(os.Stdin).Decode(&r));var mu sync.Mutex;attempts:=0;var socket net.Conn;settled:=false;started:=make(chan struct{})
 cfg:=amqp.Config{Heartbeat:0,FrameSize:131072,ChannelMax:64}
 if strings.HasPrefix(r.Scenario,"recovery-"){cfg.Recovery=&amqp.Recovery{ReconnectionConfig:&amqp.ReconnectionConfig{MaxRetryCount:2,RetryInterval:time.Millisecond}}}
 cfg.Dial=func(network,address string)(net.Conn,error){mu.Lock();attempts++;n:=attempts;mu.Unlock();if n>1&&r.Scenario=="recovery-cancel"{close(started);time.Sleep(350*time.Millisecond);mu.Lock();settled=true;mu.Unlock();return nil,errors.New("injected delayed dial failure")};s,e:=net.DialTimeout(network,address,time.Second);mu.Lock();socket=s;mu.Unlock();return s,e}
 c,e:=amqp.DialConfig(r.URI,cfg);check(e);defer c.Close();ch,e:=c.Channel();check(e)
 cn:=c.NotifyRecoveryCancel(make(chan struct{}));hn:=ch.NotifyRecoveryCancel(make(chan struct{}));notices:=c.NotifyClose(make(chan *amqp.Error,4));events:=make(chan *amqp.StateChanged,32);c.NotifyStateChange(events)
 if r.Scenario=="already"{check(c.Close())}
 if r.Scenario=="recovery-cancel"{mu.Lock();s:=socket;mu.Unlock();check(s.Close());select{case <-started:case <-time.After(5*time.Second):panic("redial not started")}}
 var deadline time.Time;if r.Scenario!="zero"{deadline=time.Now().Add(time.Duration(r.Milliseconds)*time.Millisecond)}
 begin:=time.Now();err:=c.CloseDeadline(deadline);elapsed:=time.Since(begin).Milliseconds();mu.Lock();settledAtReturn:=settled;mu.Unlock()
 if r.Scenario=="recovery-cancel"{timer:=time.NewTimer(5*time.Second);defer timer.Stop();done:=false;for !done{select{case v:=<-events:if v!=nil&&v.To==amqp.StateClosed{done=true};case <-timer.C:panic("cleanup timeout")}}}
 closeErrors:=[]any{};noticesClosed:=false;for !noticesClosed{select{case v,ok:=<-notices:if !ok{noticesClosed=true}else{closeErrors=append(closeErrors,map[string]any{"code":v.Code,"kind":kind(v)})};default:goto recorded}}
 recorded:
 out:=map[string]any{"result":kind(err),"closed":c.IsClosed(),"channelClosed":ch.IsClosed(),"connectionCancelled":cancelled(cn),"channelCancelled":cancelled(hn),"recoveryEnabled":c.IsRecoveryEnabled(),"closeErrors":closeErrors,"noticesClosed":noticesClosed,"elapsedMS":elapsed}
 if r.Scenario=="recovery-cancel"{out["dialSettledAtReturn"]=settledAtReturn}
 check(json.NewEncoder(os.Stdout).Encode(out))
}
