(logger logging.Logger, session *yamux.Session, handler http.Handler, onClose func(error), option ...Option)
| 131 | } |
| 132 | |
| 133 | func newMuxedIPC(logger logging.Logger, session *yamux.Session, handler http.Handler, onClose func(error), option ...Option) (*ipcImpl, *http.Client) { |
| 134 | // Note: Calling session.Close() needs to be done as the very last step as it shuts down all IPC! |
| 135 | |
| 136 | cfg := &cfg{shutdownTimeout: defaultShutdownTimeout} |
| 137 | for _, o := range option { |
| 138 | cfg = o(cfg) |
| 139 | } |
| 140 | server := newIpcServer(session, handler, func(err error) error { |
| 141 | if onClose != nil { |
| 142 | onClose(err) |
| 143 | } |
| 144 | return session.Close() |
| 145 | }) |
| 146 | c := createYamuxedClient(session) |
| 147 | return &ipcImpl{ |
| 148 | server: server, |
| 149 | teardown: sync.OnceValue(func() error { |
| 150 | _ = session.GoAway() |
| 151 | c.CloseIdleConnections() |
| 152 | waitForClientToDisconnect(logger, session, cfg.shutdownTimeout) |
| 153 | err := server.server.Close() |
| 154 | <-server.done |
| 155 | return errors.Join(err, server.err) |
| 156 | }), |
| 157 | }, c |
| 158 | } |
| 159 | |
| 160 | func waitForClientToDisconnect(logger logging.Logger, s *yamux.Session, t time.Duration) { |
| 161 | timeout := time.After(t) |
no test coverage detected