package iroh

import (
	
	
	
	

	
)

// ListenStreams returns a [net.Listener] view of e that accepts bidirectional
// streams as [net.Conn] values. The endpoint must already be configured with
// the ALPNs it should accept.
//
// The listener consumes e's incoming accept loop. ListenStreams returns
// [ErrEndpointAcceptLoopInUse] if [Endpoint.Accept], [Endpoint.AcceptIncoming],
// another stream listener, or [Router] already owns that loop.
//
// Closing the listener stops accepting new streams but does not close e or any
// net.Conn values already returned by [StreamListener.Accept].
func ( *Endpoint) () (*StreamListener, error) {
	if  := .acquireAcceptOwner(acceptOwnerListenStreams);  != nil {
		return nil, 
	}
	 := NewStreamListener()
	.ep = 
	.addr = net.UDPAddrFromAddrPort(.LocalAddr())
	.onClose = func() {
		.releaseAcceptOwner(acceptOwnerListenStreams)
	}
	go .run()
	return , nil
}

// NewStreamListener returns a [net.Listener] that accepts bidirectional streams
// from connections dispatched to its [StreamListener.Handler]. Register the
// handler with a [Router] to serve one ALPN as a net.Listener.
func () *StreamListener {
	,  := context.WithCancel(context.Background())
	return &StreamListener{
		ctx:     ,
		cancel:  ,
		streams: make(chan net.Conn),
		done:    make(chan struct{}),
	}
}

// StreamListener accepts bidirectional iroh streams as [net.Conn] values.
//
// Each accepted net.Conn is one bidirectional QUIC stream. Multiple accepted
// net.Conn values may come from the same peer connection. Closing an accepted
// net.Conn closes only that stream; closing the StreamListener closes any peer
// connections it has accepted but does not close the underlying endpoint. An
// accepted net.Conn also exposes RemoteID and Used0RTT methods.
type StreamListener struct {
	ep     *Endpoint
	ctx    context.Context
	cancel context.CancelFunc
	addr   net.Addr

	streams chan net.Conn
	done    chan struct{}
	onClose func()

	closeOnce sync.Once
	errMu     sync.Mutex
	err       error
}

// Accept waits for and returns the next accepted bidirectional stream.
func ( *StreamListener) () (net.Conn, error) {
	select {
	case  := <-.streams:
		return , nil
	case <-.done:
		return nil, .acceptErr()
	}
}

// Close stops accepting new streams. It does not close the underlying endpoint.
func ( *StreamListener) () error {
	.closeOnce.Do(func() {
		.setErr(net.ErrClosed)
		.cancel()
		close(.done)
		if .onClose != nil {
			.onClose()
		}
	})
	return nil
}

// Addr returns the endpoint's local UDP address.
func ( *StreamListener) () net.Addr {
	return .addr
}

// Handler returns a [ProtocolHandler] that dispatches accepted connection
// streams to l.
func ( *StreamListener) () ProtocolHandler {
	return streamListenerHandler{}
}

func ( *StreamListener) ( context.Context,  *Conn) error {
	 := make(chan struct{})
	go func() {
		.acceptStreams()
		close()
	}()
	select {
	case <-:
		return nil
	case <-.Done():
		.Close()
		<-
		return .Err()
	case <-.done:
		.Close()
		<-
		return net.ErrClosed
	}
}

func ( *StreamListener) () {
	var  sync.WaitGroup
	defer func() {
		.cancel()
		.Wait()
		.closeOnce.Do(func() {
			close(.done)
		})
	}()
	for {
		,  := .ep.accept(.ctx)
		if  != nil {
			.setErr()
			return
		}
		.Add(1)
		go func( *Conn) {
			defer .Done()
			.acceptStreams()
		}()
	}
}

type streamListenerHandler struct {
	l *StreamListener
}

func ( streamListenerHandler) ( context.Context,  *Conn) error {
	return .l.handleConn(, )
}

func ( streamListenerHandler) ( context.Context) {
	.l.Close()
}

func ( *StreamListener) ( *Conn) {
	 := newListenerConn()
	defer .doneAccepting()
	for {
		,  := .AcceptStreamConn(.ctx)
		if  != nil {
			return
		}
		.addStream()
		 = &listenerStreamConn{Conn: , owner: }
		select {
		case .streams <- :
		case <-.done:
			.Close()
			return
		case <-.ctx.Done():
			.Close()
			return
		}
	}
}

type listenerConn struct {
	conn *Conn

	mu        sync.Mutex
	active    int
	accepting bool
	closed    bool
}

func newListenerConn( *Conn) *listenerConn {
	return &listenerConn{conn: , accepting: true}
}

func ( *listenerConn) () {
	.mu.Lock()
	.active++
	.mu.Unlock()
}

func ( *listenerConn) () {
	.mu.Lock()
	if .active > 0 {
		.active--
	}
	.closeIfIdleLocked()
	.mu.Unlock()
}

func ( *listenerConn) () {
	.mu.Lock()
	.accepting = false
	.closeIfIdleLocked()
	.mu.Unlock()
}

func ( *listenerConn) () {
	if .closed || .accepting || .active != 0 {
		return
	}
	.closed = true
	.conn.Close()
}

type listenerStreamConn struct {
	net.Conn
	owner *listenerConn
	once  sync.Once
}

func ( *listenerStreamConn) () error {
	 := .Conn.Close()
	.once.Do(.owner.releaseStream)
	return 
}

func ( *listenerStreamConn) () key.EndpointID {
	return .Conn.(interface{ () key.EndpointID }).()
}

func ( *listenerStreamConn) () bool {
	return .Conn.(interface{ () bool }).()
}

func ( *StreamListener) ( error) {
	.errMu.Lock()
	defer .errMu.Unlock()
	if .err != nil {
		return
	}
	if .ctx.Err() != nil || errors.Is(, context.Canceled) {
		 = net.ErrClosed
	}
	.err = 
}

func ( *StreamListener) () error {
	.errMu.Lock()
	defer .errMu.Unlock()
	if .err != nil {
		return .err
	}
	return net.ErrClosed
}