From e3e692f0a08235826488002a1f6ba7f9900db330 Mon Sep 17 00:00:00 2001 From: David Missmann Date: Wed, 23 Sep 2026 19:54:21 +0200 Subject: [PATCH 1/8] support for file streams on XPC connections --- ios/connect.go | 4 + ios/http/http.go | 249 ++++++++++++++++++++++++++++++++++----- ios/xpc/encoding.go | 25 ++++ ios/xpc/encoding_test.go | 18 +++ ios/xpc/xpc.go | 62 +++++++++- 5 files changed, 327 insertions(+), 31 deletions(-) diff --git a/ios/connect.go b/ios/connect.go index 27fb033f3..a73e54873 100755 --- a/ios/connect.go +++ b/ios/connect.go @@ -4,6 +4,7 @@ import ( "encoding/binary" "errors" "fmt" + "io" "net" "os" "strconv" @@ -196,6 +197,9 @@ func CreateXpcConnection(h *http.HttpConnection) (*xpc.Connection, error) { if err != nil { return nil, fmt.Errorf("CreateXpcConnection: failed to create xpc connection: %w", err) } + xpcConn.SetStreamOpener(func() (io.ReadWriteCloser, error) { + return h.OpenStream() + }) return xpcConn, nil } diff --git a/ios/http/http.go b/ios/http/http.go index 194add67c..58074c653 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -4,6 +4,7 @@ import ( "bytes" "fmt" "io" + "sync" "sync/atomic" "github.com/danielpaulus/go-ios/ios/golog" @@ -20,6 +21,14 @@ const ( ServerClient = StreamId(3) ) +// defaultWindowSize is the initial HTTP/2 flow-control window (RFC 9113 6.9.2) +// that applies until the peer announces a different one. +const defaultWindowSize = 65535 + +// maxFrameSize is the largest DATA frame payload we send. It's the HTTP/2 +// default for SETTINGS_MAX_FRAME_SIZE, which every peer has to accept. +const maxFrameSize = 16384 + // HttpConnection is a wrapper around a http2.Framer that provides a simple interface to read and write http2 streams for iOS17+. type HttpConnection struct { framer *http2.Framer @@ -28,6 +37,25 @@ type HttpConnection struct { closer io.Closer csIsOpen *atomic.Bool scIsOpen *atomic.Bool + + // mu guards the flow-control state and the additional streams below. + mu sync.Mutex + // peerInitialWindow is the peer's SETTINGS_INITIAL_WINDOW_SIZE, the send + // window every newly opened stream starts with. + peerInitialWindow int64 + // connSendWindow is the connection level send window. Only writes on + // additional streams wait for it, see Stream.Write. + connSendWindow int64 + // streams holds the additional client initiated streams (5, 7, ...) that + // are used for XPC file transfers. + streams map[uint32]*streamState + nextStreamId uint32 +} + +type streamState struct { + buf bytes.Buffer + sendWindow int64 + reset bool } func (r *HttpConnection) Close() error { @@ -59,11 +87,13 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { if err != nil { return nil, fmt.Errorf("NewHttpConnection: could not read frame. %w", err) } + peerInitialWindow := int64(defaultWindowSize) if frame.Header().Type == http2.FrameSettings { settings := frame.(*http2.SettingsFrame) v, ok := settings.Value(http2.SettingInitialWindowSize) if ok { framer.SetMaxReadFrameSize(v) + peerInitialWindow = int64(v) } err := framer.WriteSettingsAck() if err != nil { @@ -80,6 +110,10 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { closer: rw, csIsOpen: &atomic.Bool{}, scIsOpen: &atomic.Bool{}, + peerInitialWindow: peerInitialWindow, + connSendWindow: defaultWindowSize, + streams: map[uint32]*streamState{}, + nextStreamId: uint32(ServerClient) + 2, }, nil } @@ -119,43 +153,102 @@ func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { if err != nil { return 0, fmt.Errorf("Write: could not write data. %w", err) } + r.mu.Lock() + r.connSendWindow -= int64(len(p)) + r.mu.Unlock() return len(p), nil } func (r *HttpConnection) readDataFrame() error { for { - f, err := r.framer.ReadFrame() + isData, err := r.processFrame() if err != nil { - return fmt.Errorf("readDataFrame: could not read frame. %w", err) - } - switch f.Header().Type { - case http2.FrameData: - d := f.(*http2.DataFrame) - switch d.StreamID { - case 1: - r.clientServerStream.Write(d.Data()) - case 3: - r.serverClientStream.Write(d.Data()) - default: - return fmt.Errorf("readDataFrame: unknown stream id %d", d.StreamID) - } + return fmt.Errorf("readDataFrame: %w", err) + } + if isData { return nil - case http2.FrameGoAway: - return fmt.Errorf("received GOAWAY") - case http2.FrameSettings: - s := f.(*http2.SettingsFrame) - if s.Flags&http2.FlagSettingsAck != http2.FlagSettingsAck { - err := r.framer.WriteSettingsAck() - if err != nil { - return fmt.Errorf("readDataFrame: could not write settings ack. %w", err) - } - } - case http2.FrameRSTStream: - r := f.(*http2.RSTStreamFrame) - return fmt.Errorf("readDataFrame: got RST frame with error code: %s", r.ErrCode.String()) + } + } +} + +// processFrame reads and handles a single frame. It reports whether the frame +// was a DATA frame so that callers waiting for stream data can re-check their +// buffers. +func (r *HttpConnection) processFrame() (bool, error) { + f, err := r.framer.ReadFrame() + if err != nil { + return false, fmt.Errorf("could not read frame. %w", err) + } + switch f.Header().Type { + case http2.FrameData: + d := f.(*http2.DataFrame) + switch d.StreamID { + case 1: + r.clientServerStream.Write(d.Data()) + case 3: + r.serverClientStream.Write(d.Data()) default: - break + r.mu.Lock() + s, ok := r.streams[d.StreamID] + if ok { + s.buf.Write(d.Data()) + } + r.mu.Unlock() + if !ok { + return false, fmt.Errorf("unknown stream id %d", d.StreamID) + } + } + return true, nil + case http2.FrameGoAway: + return false, fmt.Errorf("received GOAWAY") + case http2.FrameSettings: + s := f.(*http2.SettingsFrame) + if s.Flags&http2.FlagSettingsAck != http2.FlagSettingsAck { + if v, ok := s.Value(http2.SettingInitialWindowSize); ok { + r.updateInitialWindow(int64(v)) + } + err := r.framer.WriteSettingsAck() + if err != nil { + return false, fmt.Errorf("could not write settings ack. %w", err) + } + } + case http2.FrameWindowUpdate: + w := f.(*http2.WindowUpdateFrame) + r.mu.Lock() + if w.StreamID == uint32(InitStream) { + r.connSendWindow += int64(w.Increment) + } else if s, ok := r.streams[w.StreamID]; ok { + s.sendWindow += int64(w.Increment) + } + r.mu.Unlock() + case http2.FrameRSTStream: + rst := f.(*http2.RSTStreamFrame) + r.mu.Lock() + s, ok := r.streams[rst.StreamID] + if ok { + s.reset = true + } + r.mu.Unlock() + // The device resets file transfer streams once it received all data. + // That's no reason to fail reads on the XPC streams. + if !ok { + return false, fmt.Errorf("got RST frame with error code: %s", rst.ErrCode.String()) } + default: + break + } + return false, nil +} + +// updateInitialWindow applies a new SETTINGS_INITIAL_WINDOW_SIZE of the peer, +// which also adjusts the send windows of all open streams (RFC 9113 6.9.2). +func (r *HttpConnection) updateInitialWindow(v int64) { + r.mu.Lock() + defer r.mu.Unlock() + delta := v - r.peerInitialWindow + r.peerInitialWindow = v + for _, s := range r.streams { + s.sendWindow += delta } } @@ -200,3 +293,103 @@ func (h HttpStreamReadWriter) Write(p []byte) (n int, err error) { } return 0, fmt.Errorf("Write: unknown stream id %d", h.streamId) } + +// Stream is an additional client initiated HTTP/2 stream. RemoteXPC uses those +// for transferring the payload of file transfer objects. +// +// Writes honor the peer's flow-control windows and read frames from the +// connection while they wait for window updates. A Stream must therefore not be +// written while another goroutine reads from the same HttpConnection. +type Stream struct { + h *HttpConnection + id uint32 +} + +// OpenStream opens a new client initiated stream. +func (r *HttpConnection) OpenStream() (*Stream, error) { + r.mu.Lock() + id := r.nextStreamId + r.nextStreamId += 2 + r.streams[id] = &streamState{sendWindow: r.peerInitialWindow} + r.mu.Unlock() + + err := r.framer.WriteHeaders(http2.HeadersFrameParam{ + StreamID: id, + EndHeaders: true, + }) + if err != nil { + return nil, fmt.Errorf("OpenStream: could not send headers for stream %d. %w", id, err) + } + return &Stream{h: r, id: id}, nil +} + +// Read blocks until len(p) bytes were received on this stream +func (s *Stream) Read(p []byte) (int, error) { + for { + s.h.mu.Lock() + st := s.h.streams[s.id] + if st.buf.Len() >= len(p) { + n, err := st.buf.Read(p) + s.h.mu.Unlock() + return n, err + } + reset := st.reset + s.h.mu.Unlock() + if reset { + return 0, fmt.Errorf("Read: stream %d was reset by the peer", s.id) + } + if _, err := s.h.processFrame(); err != nil { + return 0, fmt.Errorf("Read: %w", err) + } + } +} + +// Write sends p as DATA frames on this stream, waiting for the peer to open its +// flow-control windows whenever they are exhausted. +func (s *Stream) Write(p []byte) (int, error) { + written := 0 + for written < len(p) { + n, err := s.sendWindow(len(p) - written) + if err != nil { + return written, fmt.Errorf("Write: %w", err) + } + if n == 0 { + if _, err := s.h.processFrame(); err != nil { + return written, fmt.Errorf("Write: failed waiting for window update. %w", err) + } + continue + } + if err := s.h.framer.WriteData(s.id, false, p[written:written+n]); err != nil { + return written, fmt.Errorf("Write: could not write data on stream %d. %w", s.id, err) + } + written += n + } + return written, nil +} + +// sendWindow reserves up to want bytes of the stream and connection windows and +// returns how many bytes may be sent right now. +func (s *Stream) sendWindow(want int) (int, error) { + s.h.mu.Lock() + defer s.h.mu.Unlock() + st := s.h.streams[s.id] + if st.reset { + return 0, fmt.Errorf("stream %d was reset by the peer", s.id) + } + n := int64(min(want, maxFrameSize)) + n = min(n, st.sendWindow, s.h.connSendWindow) + if n <= 0 { + return 0, nil + } + st.sendWindow -= n + s.h.connSendWindow -= n + return int(n), nil +} + +// Close half-closes the stream by sending an empty DATA frame with END_STREAM +func (s *Stream) Close() error { + if err := s.h.framer.WriteData(s.id, true, nil); err != nil { + return fmt.Errorf("Close: could not end stream %d. %w", s.id, err) + } + return nil +} diff --git a/ios/xpc/encoding.go b/ios/xpc/encoding.go index f7e2dadaf..71035eb94 100644 --- a/ios/xpc/encoding.go +++ b/ios/xpc/encoding.go @@ -44,6 +44,7 @@ const ( HeartbeatRequestFlag = uint32(0x00010000) HeartbeatReplyFlag = uint32(0x00020000) FileOpenFlag = uint32(0x00100000) + FileOpenReplyFlag = uint32(0x00200000) InitHandshakeFlag = uint32(0x00400000) ) @@ -156,6 +157,7 @@ func decodeWrapper(r io.Reader) (Message, error) { if h.BodyLen == 0 { return Message{ Flags: h.Flags, + Id: h.MsgId, }, nil } body, err := decodeBody(r, h) @@ -165,6 +167,7 @@ func decodeWrapper(r io.Reader) (Message, error) { return Message{ Flags: h.Flags, Body: body, + Id: h.MsgId, }, nil } @@ -551,6 +554,10 @@ func encodeObject(w io.Writer, e interface{}) error { if err := encodeDictionary(w, e.(map[string]interface{})); err != nil { return err } + case FileTransfer: + if err := encodeFileTransfer(w, t); err != nil { + return err + } default: return fmt.Errorf("can not encode type %v", t) } @@ -569,6 +576,24 @@ func encodeUuid(w io.Writer, u uuid.UUID) error { return nil } +// encodeFileTransfer writes a file transfer object. The payload itself is not +// part of the message, it gets sent on a separate stream that is opened with a +// FileOpenFlag message carrying the same MsgId. +func encodeFileTransfer(w io.Writer, f FileTransfer) error { + header := struct { + t xpcType + msgId uint64 + }{fileTransferType, f.MsgId} + if err := binary.Write(w, binary.LittleEndian, header); err != nil { + return fmt.Errorf("encodeFileTransfer: failed to write header: %w", err) + } + // the transfer length is always stored in a property 's' + if err := encodeDictionary(w, map[string]interface{}{"s": f.TransferSize}); err != nil { + return fmt.Errorf("encodeFileTransfer: failed to write transfer size: %w", err) + } + return nil +} + func encodeArray(w io.Writer, slice []interface{}) error { buf := bytes.NewBuffer(nil) for i, e := range slice { diff --git a/ios/xpc/encoding_test.go b/ios/xpc/encoding_test.go index a04d72f92..92296b882 100644 --- a/ios/xpc/encoding_test.go +++ b/ios/xpc/encoding_test.go @@ -3,6 +3,7 @@ package xpc import ( "bytes" "encoding/base64" + "encoding/hex" "github.com/google/uuid" "github.com/stretchr/testify/assert" "os" @@ -62,6 +63,7 @@ func TestDictionary(t *testing.T) { }, "CoreDevice.invocationIdentifier": "62419FC1-5ABF-4D96-BCA8-7A5F6F9A69EE", }, + Id: 1, }, res) } @@ -126,6 +128,13 @@ func TestEncodeDecode(t *testing.T) { }, expectedFlags: AlwaysSetFlag | DataFlag, }, + { + name: "encode file transfer", + input: map[string]interface{}{ + "image": FileTransfer{MsgId: 13, TransferSize: 16648704}, + }, + expectedFlags: AlwaysSetFlag | DataFlag, + }, { name: "encode uuid", input: map[string]interface{}{ @@ -164,3 +173,12 @@ func TestEncodeDecode(t *testing.T) { }) } } + +// TestEncodeFileTransferWireFormat checks the encoding against the bytes a Mac +// sends for the 'image' argument of the cryptexd 'install' routine. +func TestEncodeFileTransferWireFormat(t *testing.T) { + buf := bytes.NewBuffer(nil) + err := encodeObject(buf, FileTransfer{MsgId: 13, TransferSize: 16648704}) + assert.NoError(t, err) + assert.Equal(t, "00a001000d0000000000000000f0000014000000010000007300000000400000000afe0000000000", hex.EncodeToString(buf.Bytes())) +} diff --git a/ios/xpc/xpc.go b/ios/xpc/xpc.go index 3ff2cb150..ee4d562d6 100644 --- a/ios/xpc/xpc.go +++ b/ios/xpc/xpc.go @@ -5,19 +5,21 @@ package xpc import ( "fmt" "io" - - "golang.org/x/net/http2" ) // Connection represents a http2 based connection to an XPC service on an iOS17 device. type Connection struct { connectionCloser io.Closer - framer *http2.Framer msgId uint64 clientServer io.ReadWriter serverClient io.ReadWriter + openStream StreamOpener } +// StreamOpener opens an additional stream on the underlying connection. RemoteXPC +// sends the payload of each FileTransfer object on a stream of its own. +type StreamOpener func() (io.ReadWriteCloser, error) + // New creates a new connection to an XPC service on an iOS17 device. func New(clientServer io.ReadWriter, serverClient io.ReadWriter, closer io.Closer) (*Connection, error) { return &Connection{ @@ -66,6 +68,60 @@ func (c *Connection) Send(data map[string]interface{}, flags ...uint32) error { return EncodeMessage(c.clientServer, msg) } +// SetStreamOpener enables OpenFileTransfer on this connection +func (c *Connection) SetStreamOpener(o StreamOpener) { + c.openStream = o +} + +// FileTransferStream carries the payload of a FileTransfer object +type FileTransferStream struct { + rwc io.ReadWriteCloser + id uint64 +} + +// OpenFileTransfer opens a new stream for the payload of the FileTransfer object with +// the given id. The peer accepts it only after it received the message that +// references the FileTransfer, so call WaitAccepted after sending that message. +func (c *Connection) OpenFileTransfer(id uint64) (*FileTransferStream, error) { + if c.openStream == nil { + return nil, fmt.Errorf("OpenFileTransfer: connection does not support file transfers") + } + rwc, err := c.openStream() + if err != nil { + return nil, fmt.Errorf("OpenFileTransfer: %w", err) + } + err = EncodeMessage(rwc, Message{ + Flags: AlwaysSetFlag | FileOpenFlag, + Id: id, + }) + if err != nil { + return nil, fmt.Errorf("OpenFileTransfer: failed to send file open message: %w", err) + } + return &FileTransferStream{rwc: rwc, id: id}, nil +} + +// WaitAccepted blocks until the peer is ready to receive the payload +func (f *FileTransferStream) WaitAccepted() error { + msg, err := DecodeMessage(f.rwc) + if err != nil { + return fmt.Errorf("WaitAccepted: %w", err) + } + if msg.Flags&FileOpenReplyFlag == 0 || msg.Id != f.id { + return fmt.Errorf("WaitAccepted: unexpected reply with flags 0x%x and id %d for file transfer %d", msg.Flags, msg.Id, f.id) + } + return nil +} + +// Write sends payload data +func (f *FileTransferStream) Write(p []byte) (int, error) { + return f.rwc.Write(p) +} + +// Close marks the end of the payload +func (f *FileTransferStream) Close() error { + return f.rwc.Close() +} + func (c *Connection) Close() error { return c.connectionCloser.Close() } From 41112df17f510e00c5da845e6c675a67e55c6a5b Mon Sep 17 00:00:00 2001 From: David Missmann Date: Fri, 25 Sep 2026 13:30:12 +0200 Subject: [PATCH 2/8] use the same stream implementation for all HTTP streams --- ios/connect.go | 18 +++---- ios/http/http.go | 121 +++++++---------------------------------------- 2 files changed, 26 insertions(+), 113 deletions(-) diff --git a/ios/connect.go b/ios/connect.go index a73e54873..ebe1b65b5 100755 --- a/ios/connect.go +++ b/ios/connect.go @@ -185,14 +185,19 @@ func ConnectToServiceTunnelIface(device DeviceEntry, serviceName string) (Device } func CreateXpcConnection(h *http.HttpConnection) (*xpc.Connection, error) { - err := initializeXpcConnection(h) + clientServerChannel, err := h.OpenStream() + if err != nil { + return nil, fmt.Errorf("CreateXpcConnection: failed to open stream: %w", err) + } + serverClientChannel, err := h.OpenStream() + if err != nil { + return nil, fmt.Errorf("CreateXpcConnection: failed to open stream: %w", err) + } + err = initializeXpcConnection(clientServerChannel, serverClientChannel) if err != nil { return nil, fmt.Errorf("CreateXpcConnection: failed to initialize xpc connection: %w", err) } - clientServerChannel := http.NewStreamReadWriter(h, http.ClientServer) - serverClientChannel := http.NewStreamReadWriter(h, http.ServerClient) - xpcConn, err := xpc.New(clientServerChannel, serverClientChannel, h) if err != nil { return nil, fmt.Errorf("CreateXpcConnection: failed to create xpc connection: %w", err) @@ -251,10 +256,7 @@ func ConnectLockdownWithSession(device DeviceEntry) (*LockDownConnection, error) return lockdownConnection, nil } -func initializeXpcConnection(h *http.HttpConnection) error { - csWriter := http.NewStreamReadWriter(h, http.ClientServer) - ssWriter := http.NewStreamReadWriter(h, http.ServerClient) - +func initializeXpcConnection(csWriter, ssWriter io.ReadWriter) error { err := xpc.EncodeMessage(csWriter, xpc.Message{ Flags: xpc.AlwaysSetFlag, Body: map[string]interface{}{}, diff --git a/ios/http/http.go b/ios/http/http.go index 58074c653..b77612d65 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -5,7 +5,6 @@ import ( "fmt" "io" "sync" - "sync/atomic" "github.com/danielpaulus/go-ios/ios/golog" "golang.org/x/net/http2" @@ -31,12 +30,8 @@ const maxFrameSize = 16384 // HttpConnection is a wrapper around a http2.Framer that provides a simple interface to read and write http2 streams for iOS17+. type HttpConnection struct { - framer *http2.Framer - clientServerStream *bytes.Buffer - serverClientStream *bytes.Buffer - closer io.Closer - csIsOpen *atomic.Bool - scIsOpen *atomic.Bool + framer *http2.Framer + closer io.Closer // mu guards the flow-control state and the additional streams below. mu sync.Mutex @@ -104,50 +99,15 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { } return &HttpConnection{ - framer: framer, - clientServerStream: bytes.NewBuffer(nil), - serverClientStream: bytes.NewBuffer(nil), - closer: rw, - csIsOpen: &atomic.Bool{}, - scIsOpen: &atomic.Bool{}, - peerInitialWindow: peerInitialWindow, - connSendWindow: defaultWindowSize, - streams: map[uint32]*streamState{}, - nextStreamId: uint32(ServerClient) + 2, + framer: framer, + closer: rw, + peerInitialWindow: peerInitialWindow, + connSendWindow: defaultWindowSize, + streams: map[uint32]*streamState{}, + nextStreamId: uint32(1), }, nil } -func (r *HttpConnection) ReadClientServerStream(p []byte) (int, error) { - for r.clientServerStream.Len() < len(p) { - err := r.readDataFrame() - if err != nil { - return 0, fmt.Errorf("ReadClientServerStream: %w", err) - } - } - return r.clientServerStream.Read(p) -} - -func (r *HttpConnection) WriteClientServerStream(p []byte) (int, error) { - return r.write(p, uint32(ClientServer), r.csIsOpen) -} - -func (r *HttpConnection) WriteServerClientStream(p []byte) (int, error) { - return r.write(p, uint32(ServerClient), r.scIsOpen) -} - -func (r *HttpConnection) write(p []byte, stream uint32, isOpen *atomic.Bool) (int, error) { - if isOpen.CompareAndSwap(false, true) { - err := r.framer.WriteHeaders(http2.HeadersFrameParam{ - StreamID: stream, - EndHeaders: true, - }) - if err != nil { - return 0, fmt.Errorf("write: could not send headers. %w", err) - } - } - return r.Write(p, stream) -} - func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { err := r.framer.WriteData(streamId, false, p) if err != nil { @@ -182,21 +142,14 @@ func (r *HttpConnection) processFrame() (bool, error) { switch f.Header().Type { case http2.FrameData: d := f.(*http2.DataFrame) - switch d.StreamID { - case 1: - r.clientServerStream.Write(d.Data()) - case 3: - r.serverClientStream.Write(d.Data()) - default: - r.mu.Lock() - s, ok := r.streams[d.StreamID] - if ok { - s.buf.Write(d.Data()) - } - r.mu.Unlock() - if !ok { - return false, fmt.Errorf("unknown stream id %d", d.StreamID) - } + r.mu.Lock() + s, ok := r.streams[d.StreamID] + if ok { + s.buf.Write(d.Data()) + } + r.mu.Unlock() + if !ok { + return false, fmt.Errorf("unknown stream id %d", d.StreamID) } return true, nil case http2.FrameGoAway: @@ -252,48 +205,6 @@ func (r *HttpConnection) updateInitialWindow(v int64) { } } -func (r *HttpConnection) ReadServerClientStream(p []byte) (int, error) { - for r.serverClientStream.Len() < len(p) { - err := r.readDataFrame() - if err != nil { - return 0, err - } - } - return r.serverClientStream.Read(p) -} - -type HttpStreamReadWriter struct { - h *HttpConnection - streamId uint32 -} - -func NewStreamReadWriter(h *HttpConnection, streamId StreamId) HttpStreamReadWriter { - return HttpStreamReadWriter{ - h: h, - streamId: uint32(streamId), - } -} - -func (h HttpStreamReadWriter) Read(p []byte) (n int, err error) { - if h.streamId == 1 { - return h.h.ReadClientServerStream(p) - } - if h.streamId == 3 { - return h.h.ReadServerClientStream(p) - } - return 0, fmt.Errorf("Read: unknown stream id %d", h.streamId) -} - -func (h HttpStreamReadWriter) Write(p []byte) (n int, err error) { - if h.streamId == 1 { - return h.h.WriteClientServerStream(p) - } - if h.streamId == 3 { - return h.h.WriteServerClientStream(p) - } - return 0, fmt.Errorf("Write: unknown stream id %d", h.streamId) -} - // Stream is an additional client initiated HTTP/2 stream. RemoteXPC uses those // for transferring the payload of file transfer objects. // From 0404c8fe2be1181edb215497eb36bee6ebb88dfe Mon Sep 17 00:00:00 2001 From: David Missmann Date: Fri, 25 Sep 2026 13:44:57 +0200 Subject: [PATCH 3/8] remove setter for open stream --- ios/connect.go | 8 ++++---- ios/xpc/xpc.go | 8 ++------ 2 files changed, 6 insertions(+), 10 deletions(-) diff --git a/ios/connect.go b/ios/connect.go index ebe1b65b5..7ed1511b5 100755 --- a/ios/connect.go +++ b/ios/connect.go @@ -198,13 +198,13 @@ func CreateXpcConnection(h *http.HttpConnection) (*xpc.Connection, error) { return nil, fmt.Errorf("CreateXpcConnection: failed to initialize xpc connection: %w", err) } - xpcConn, err := xpc.New(clientServerChannel, serverClientChannel, h) + openStream := func() (io.ReadWriteCloser, error) { + return h.OpenStream() + } + xpcConn, err := xpc.New(clientServerChannel, serverClientChannel, openStream, h) if err != nil { return nil, fmt.Errorf("CreateXpcConnection: failed to create xpc connection: %w", err) } - xpcConn.SetStreamOpener(func() (io.ReadWriteCloser, error) { - return h.OpenStream() - }) return xpcConn, nil } diff --git a/ios/xpc/xpc.go b/ios/xpc/xpc.go index ee4d562d6..9d0009ed9 100644 --- a/ios/xpc/xpc.go +++ b/ios/xpc/xpc.go @@ -21,12 +21,13 @@ type Connection struct { type StreamOpener func() (io.ReadWriteCloser, error) // New creates a new connection to an XPC service on an iOS17 device. -func New(clientServer io.ReadWriter, serverClient io.ReadWriter, closer io.Closer) (*Connection, error) { +func New(clientServer io.ReadWriter, serverClient io.ReadWriter, openStream StreamOpener, closer io.Closer) (*Connection, error) { return &Connection{ connectionCloser: closer, msgId: 1, clientServer: clientServer, serverClient: serverClient, + openStream: openStream, }, nil } @@ -68,11 +69,6 @@ func (c *Connection) Send(data map[string]interface{}, flags ...uint32) error { return EncodeMessage(c.clientServer, msg) } -// SetStreamOpener enables OpenFileTransfer on this connection -func (c *Connection) SetStreamOpener(o StreamOpener) { - c.openStream = o -} - // FileTransferStream carries the payload of a FileTransfer object type FileTransferStream struct { rwc io.ReadWriteCloser From 1a4deb705a37e6bbf38f5fe6347c7273f4ea3f2c Mon Sep 17 00:00:00 2001 From: David Missmann Date: Mon, 28 Sep 2026 12:34:17 +0200 Subject: [PATCH 4/8] remove unused method --- ios/http/http.go | 37 ++++++++++++------------------------- 1 file changed, 12 insertions(+), 25 deletions(-) diff --git a/ios/http/http.go b/ios/http/http.go index b77612d65..0d8a2ab4c 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -119,25 +119,13 @@ func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { return len(p), nil } -func (r *HttpConnection) readDataFrame() error { - for { - isData, err := r.processFrame() - if err != nil { - return fmt.Errorf("readDataFrame: %w", err) - } - if isData { - return nil - } - } -} - -// processFrame reads and handles a single frame. It reports whether the frame -// was a DATA frame so that callers waiting for stream data can re-check their -// buffers. -func (r *HttpConnection) processFrame() (bool, error) { +// processFrame reads and handles a single frame. Callers waiting for stream +// data or for flow-control windows to open call it repeatedly until their +// condition is met. +func (r *HttpConnection) processFrame() error { f, err := r.framer.ReadFrame() if err != nil { - return false, fmt.Errorf("could not read frame. %w", err) + return fmt.Errorf("could not read frame. %w", err) } switch f.Header().Type { case http2.FrameData: @@ -149,11 +137,10 @@ func (r *HttpConnection) processFrame() (bool, error) { } r.mu.Unlock() if !ok { - return false, fmt.Errorf("unknown stream id %d", d.StreamID) + return fmt.Errorf("unknown stream id %d", d.StreamID) } - return true, nil case http2.FrameGoAway: - return false, fmt.Errorf("received GOAWAY") + return fmt.Errorf("received GOAWAY") case http2.FrameSettings: s := f.(*http2.SettingsFrame) if s.Flags&http2.FlagSettingsAck != http2.FlagSettingsAck { @@ -162,7 +149,7 @@ func (r *HttpConnection) processFrame() (bool, error) { } err := r.framer.WriteSettingsAck() if err != nil { - return false, fmt.Errorf("could not write settings ack. %w", err) + return fmt.Errorf("could not write settings ack. %w", err) } } case http2.FrameWindowUpdate: @@ -185,12 +172,12 @@ func (r *HttpConnection) processFrame() (bool, error) { // The device resets file transfer streams once it received all data. // That's no reason to fail reads on the XPC streams. if !ok { - return false, fmt.Errorf("got RST frame with error code: %s", rst.ErrCode.String()) + return fmt.Errorf("got RST frame with error code: %s", rst.ErrCode.String()) } default: break } - return false, nil + return nil } // updateInitialWindow applies a new SETTINGS_INITIAL_WINDOW_SIZE of the peer, @@ -249,7 +236,7 @@ func (s *Stream) Read(p []byte) (int, error) { if reset { return 0, fmt.Errorf("Read: stream %d was reset by the peer", s.id) } - if _, err := s.h.processFrame(); err != nil { + if err := s.h.processFrame(); err != nil { return 0, fmt.Errorf("Read: %w", err) } } @@ -265,7 +252,7 @@ func (s *Stream) Write(p []byte) (int, error) { return written, fmt.Errorf("Write: %w", err) } if n == 0 { - if _, err := s.h.processFrame(); err != nil { + if err := s.h.processFrame(); err != nil { return written, fmt.Errorf("Write: failed waiting for window update. %w", err) } continue From 0de7aa7c41085277bb6c30159d4812feb4f689ec Mon Sep 17 00:00:00 2001 From: David Missmann Date: Mon, 28 Sep 2026 12:36:09 +0200 Subject: [PATCH 5/8] add comment about stream IDs --- ios/http/http.go | 1 + 1 file changed, 1 insertion(+) diff --git a/ios/http/http.go b/ios/http/http.go index 0d8a2ab4c..1326e1edb 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -207,6 +207,7 @@ type Stream struct { func (r *HttpConnection) OpenStream() (*Stream, error) { r.mu.Lock() id := r.nextStreamId + // Client controlled streams are always odd numbered (RFC 9113 5.1.1) r.nextStreamId += 2 r.streams[id] = &streamState{sendWindow: r.peerInitialWindow} r.mu.Unlock() From 8872d91d17bf0ea29cc79acf38d0f3b4b59e435a Mon Sep 17 00:00:00 2001 From: David Missmann Date: Tue, 29 Sep 2026 12:14:31 +0200 Subject: [PATCH 6/8] lock access to http2.Framer --- ios/http/http.go | 66 ++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 58 insertions(+), 8 deletions(-) diff --git a/ios/http/http.go b/ios/http/http.go index 1326e1edb..36be0f38b 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -30,9 +30,18 @@ const maxFrameSize = 16384 // HttpConnection is a wrapper around a http2.Framer that provides a simple interface to read and write http2 streams for iOS17+. type HttpConnection struct { - framer *http2.Framer closer io.Closer + framer *http2.Framer + + // framerReadMu guards reading frames. A http2.Framer is not safe for concurrent + // reads, and the payload of a frame is only valid until the next read, so a + // frame has to be dispatched to its stream before the next one is read. + framerReadMu sync.Mutex + // framerWriteMu guards writing frames, a http2.Framer encodes every frame into + // the same buffer. + framerWriteMu sync.Mutex + // mu guards the flow-control state and the additional streams below. mu sync.Mutex // peerInitialWindow is the peer's SETTINGS_INITIAL_WINDOW_SIZE, the send @@ -109,7 +118,9 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { } func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { + r.framerWriteMu.Lock() err := r.framer.WriteData(streamId, false, p) + r.framerWriteMu.Unlock() if err != nil { return 0, fmt.Errorf("Write: could not write data. %w", err) } @@ -122,7 +133,18 @@ func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { // processFrame reads and handles a single frame. Callers waiting for stream // data or for flow-control windows to open call it repeatedly until their // condition is met. -func (r *HttpConnection) processFrame() error { +// +// Only one goroutine reads from the connection at a time, and it dispatches the +// frames of all streams. ready is therefore evaluated again once this goroutine +// owns the read side: another goroutine may have delivered what the caller is +// waiting for in the meantime, and reading one more frame would block until the +// peer happens to send another one. +func (r *HttpConnection) processFrame(ready func() bool) error { + r.framerReadMu.Lock() + defer r.framerReadMu.Unlock() + if ready() { + return nil + } f, err := r.framer.ReadFrame() if err != nil { return fmt.Errorf("could not read frame. %w", err) @@ -147,7 +169,9 @@ func (r *HttpConnection) processFrame() error { if v, ok := s.Value(http2.SettingInitialWindowSize); ok { r.updateInitialWindow(int64(v)) } + r.framerWriteMu.Lock() err := r.framer.WriteSettingsAck() + r.framerWriteMu.Unlock() if err != nil { return fmt.Errorf("could not write settings ack. %w", err) } @@ -196,8 +220,8 @@ func (r *HttpConnection) updateInitialWindow(v int64) { // for transferring the payload of file transfer objects. // // Writes honor the peer's flow-control windows and read frames from the -// connection while they wait for window updates. A Stream must therefore not be -// written while another goroutine reads from the same HttpConnection. +// connection while they wait for window updates. Every stream of a connection +// can be read and written from its own goroutine. type Stream struct { h *HttpConnection id uint32 @@ -212,10 +236,12 @@ func (r *HttpConnection) OpenStream() (*Stream, error) { r.streams[id] = &streamState{sendWindow: r.peerInitialWindow} r.mu.Unlock() + r.framerWriteMu.Lock() err := r.framer.WriteHeaders(http2.HeadersFrameParam{ StreamID: id, EndHeaders: true, }) + r.framerWriteMu.Unlock() if err != nil { return nil, fmt.Errorf("OpenStream: could not send headers for stream %d. %w", id, err) } @@ -237,12 +263,21 @@ func (s *Stream) Read(p []byte) (int, error) { if reset { return 0, fmt.Errorf("Read: stream %d was reset by the peer", s.id) } - if err := s.h.processFrame(); err != nil { + if err := s.h.processFrame(func() bool { return s.readable(len(p)) }); err != nil { return 0, fmt.Errorf("Read: %w", err) } } } +// readable reports whether a read of n bytes can complete, either because the +// data arrived or because the peer reset the stream. +func (s *Stream) readable(n int) bool { + s.h.mu.Lock() + defer s.h.mu.Unlock() + st := s.h.streams[s.id] + return st.reset || st.buf.Len() >= n +} + // Write sends p as DATA frames on this stream, waiting for the peer to open its // flow-control windows whenever they are exhausted. func (s *Stream) Write(p []byte) (int, error) { @@ -253,12 +288,15 @@ func (s *Stream) Write(p []byte) (int, error) { return written, fmt.Errorf("Write: %w", err) } if n == 0 { - if err := s.h.processFrame(); err != nil { + if err := s.h.processFrame(s.writable); err != nil { return written, fmt.Errorf("Write: failed waiting for window update. %w", err) } continue } - if err := s.h.framer.WriteData(s.id, false, p[written:written+n]); err != nil { + s.h.framerWriteMu.Lock() + err = s.h.framer.WriteData(s.id, false, p[written:written+n]) + s.h.framerWriteMu.Unlock() + if err != nil { return written, fmt.Errorf("Write: could not write data on stream %d. %w", s.id, err) } written += n @@ -266,6 +304,15 @@ func (s *Stream) Write(p []byte) (int, error) { return written, nil } +// writable reports whether a write can make progress, either because both send +// windows have room or because the peer reset the stream. +func (s *Stream) writable() bool { + s.h.mu.Lock() + defer s.h.mu.Unlock() + st := s.h.streams[s.id] + return st.reset || (st.sendWindow > 0 && s.h.connSendWindow > 0) +} + // sendWindow reserves up to want bytes of the stream and connection windows and // returns how many bytes may be sent right now. func (s *Stream) sendWindow(want int) (int, error) { @@ -287,7 +334,10 @@ func (s *Stream) sendWindow(want int) (int, error) { // Close half-closes the stream by sending an empty DATA frame with END_STREAM func (s *Stream) Close() error { - if err := s.h.framer.WriteData(s.id, true, nil); err != nil { + s.h.framerWriteMu.Lock() + err := s.h.framer.WriteData(s.id, true, nil) + s.h.framerWriteMu.Unlock() + if err != nil { return fmt.Errorf("Close: could not end stream %d. %w", s.id, err) } return nil From 91549bca4a4924e649988302847502b28e2168d4 Mon Sep 17 00:00:00 2001 From: David Missmann Date: Tue, 29 Sep 2026 14:06:02 +0200 Subject: [PATCH 7/8] send WINDOW_UPDATE messages for received data frames --- ios/http/http.go | 67 +++++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 64 insertions(+), 3 deletions(-) diff --git a/ios/http/http.go b/ios/http/http.go index 36be0f38b..0f5f25812 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -24,6 +24,15 @@ const ( // that applies until the peer announces a different one. const defaultWindowSize = 65535 +// recvWindowSize is the flow-control window we grant the peer, for the +// connection as well as for every stream. It is announced as +// SETTINGS_INITIAL_WINDOW_SIZE and replenished with WINDOW_UPDATE frames. +const recvWindowSize = 1048576 + +// windowUpdateThreshold is how much of a receive window may be used up before +// it is replenished, so that not every DATA frame needs a WINDOW_UPDATE. +const windowUpdateThreshold = recvWindowSize / 2 + // maxFrameSize is the largest DATA frame payload we send. It's the HTTP/2 // default for SETTINGS_MAX_FRAME_SIZE, which every peer has to accept. const maxFrameSize = 16384 @@ -50,6 +59,9 @@ type HttpConnection struct { // connSendWindow is the connection level send window. Only writes on // additional streams wait for it, see Stream.Write. connSendWindow int64 + // connRecvUnacked counts the bytes received on the connection that have not + // been granted back to the peer with a WINDOW_UPDATE yet. + connRecvUnacked int64 // streams holds the additional client initiated streams (5, 7, ...) that // are used for XPC file transfers. streams map[uint32]*streamState @@ -59,7 +71,10 @@ type HttpConnection struct { type streamState struct { buf bytes.Buffer sendWindow int64 - reset bool + // recvUnacked counts the bytes received on this stream that have not been + // granted back to the peer with a WINDOW_UPDATE yet. + recvUnacked int64 + reset bool } func (r *HttpConnection) Close() error { @@ -76,13 +91,15 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { err = framer.WriteSettings( http2.Setting{ID: http2.SettingMaxConcurrentStreams, Val: 100}, - http2.Setting{ID: http2.SettingInitialWindowSize, Val: 1048576}, + http2.Setting{ID: http2.SettingInitialWindowSize, Val: recvWindowSize}, ) if err != nil { return nil, fmt.Errorf("NewHttpConnection: could not write settings. %w", err) } - err = framer.WriteWindowUpdate(uint32(InitStream), 983041) + // SETTINGS_INITIAL_WINDOW_SIZE doesn't apply to the connection window, so it + // is raised to the same size explicitly (RFC 9113 6.9.2) + err = framer.WriteWindowUpdate(uint32(InitStream), recvWindowSize-defaultWindowSize) if err != nil { return nil, fmt.Errorf("NewHttpConnection: could not write window update. %w", err) } @@ -152,15 +169,32 @@ func (r *HttpConnection) processFrame(ready func() bool) error { switch f.Header().Type { case http2.FrameData: d := f.(*http2.DataFrame) + // the whole frame payload counts against the receive windows, the + // padding included (RFC 9113 6.9.1) + size := int64(d.Header().Length) r.mu.Lock() s, ok := r.streams[d.StreamID] + streamIncrement := int64(0) if ok { s.buf.Write(d.Data()) + // a stream that the peer is done with doesn't need its window back + if !s.reset && !d.StreamEnded() { + streamIncrement = ackReceived(&s.recvUnacked, size) + } } + connIncrement := ackReceived(&r.connRecvUnacked, size) r.mu.Unlock() if !ok { return fmt.Errorf("unknown stream id %d", d.StreamID) } + // the received data is buffered without a bound, so the windows are + // replenished right away instead of when a reader consumes the data + if err := r.writeWindowUpdate(uint32(InitStream), connIncrement); err != nil { + return err + } + if err := r.writeWindowUpdate(d.StreamID, streamIncrement); err != nil { + return err + } case http2.FrameGoAway: return fmt.Errorf("received GOAWAY") case http2.FrameSettings: @@ -216,6 +250,33 @@ func (r *HttpConnection) updateInitialWindow(v int64) { } } +// ackReceived adds the n received bytes to unacked and returns the increment to +// grant back with a WINDOW_UPDATE frame, or 0 while the window still has room. +func ackReceived(unacked *int64, n int64) int64 { + *unacked += n + if *unacked < windowUpdateThreshold { + return 0 + } + increment := *unacked + *unacked = 0 + return increment +} + +// writeWindowUpdate grants increment bytes of receive window back to the peer. +// The WINDOW_UPDATE is only sent when increment > 0, otherwise this is a no-op. +func (r *HttpConnection) writeWindowUpdate(streamId uint32, increment int64) error { + if increment <= 0 { + return nil + } + r.framerWriteMu.Lock() + err := r.framer.WriteWindowUpdate(streamId, uint32(increment)) + r.framerWriteMu.Unlock() + if err != nil { + return fmt.Errorf("could not write window update for stream %d. %w", streamId, err) + } + return nil +} + // Stream is an additional client initiated HTTP/2 stream. RemoteXPC uses those // for transferring the payload of file transfer objects. // From 7224ce9e091c23f0040d014767b7ca105e4da240 Mon Sep 17 00:00:00 2001 From: David Missmann Date: Tue, 29 Sep 2026 14:16:28 +0200 Subject: [PATCH 8/8] remove unused HttpConnection.Write method --- ios/http/http.go | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/ios/http/http.go b/ios/http/http.go index 0f5f25812..0fbed0878 100644 --- a/ios/http/http.go +++ b/ios/http/http.go @@ -134,19 +134,6 @@ func NewHttpConnection(rw io.ReadWriteCloser) (*HttpConnection, error) { }, nil } -func (r *HttpConnection) Write(p []byte, streamId uint32) (int, error) { - r.framerWriteMu.Lock() - err := r.framer.WriteData(streamId, false, p) - r.framerWriteMu.Unlock() - if err != nil { - return 0, fmt.Errorf("Write: could not write data. %w", err) - } - r.mu.Lock() - r.connSendWindow -= int64(len(p)) - r.mu.Unlock() - return len(p), nil -} - // processFrame reads and handles a single frame. Callers waiting for stream // data or for flow-control windows to open call it repeatedly until their // condition is met.