aboutsummaryrefslogtreecommitdiffhomepage
path: root/helper/proc/pipe.go
blob: 838c4ae4c50df9b0ae03e8a708da834f01570369 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
package proc

import (
	"context"
	"errors"
	"io"
	"os"
	"runtime"
)

// NewWriterTo returns a [File] that receives content from wt on fulfillment.
func NewWriterTo(wt io.WriterTo) File { return &writeToFile{wt: wt} }

// writeToFile exports the read end of a pipe with data written by an [io.WriterTo].
type writeToFile struct {
	wt io.WriterTo
	BaseFile
}

func (f *writeToFile) ErrCount() int { return 3 }
func (f *writeToFile) Fulfill(ctx context.Context, dispatchErr func(error)) error {
	r, w, err := os.Pipe()
	if err != nil {
		return err
	}
	f.Set(r)

	done := make(chan struct{})
	go func() {
		_, err = f.wt.WriteTo(w)
		dispatchErr(err)
		dispatchErr(w.Close())
		close(done)
		runtime.KeepAlive(r)
	}()
	go func() {
		select {
		case <-done:
			dispatchErr(nil)
		case <-ctx.Done():
			dispatchErr(w.Close()) // this aborts WriteTo with file already closed
			runtime.KeepAlive(r)
		}
	}()

	return nil
}

// NewStat returns a [File] implementing the behaviour
// of the receiving end of xdg-dbus-proxy stat fd.
func NewStat(s *io.Closer) File { return &statFile{s: s} }

var (
	ErrStatFault = errors.New("generic stat fd fault")
	ErrStatRead  = errors.New("unexpected stat behaviour")
)

// statFile implements xdg-dbus-proxy stat fd behaviour.
type statFile struct {
	s *io.Closer
	BaseFile
}

func (f *statFile) ErrCount() int { return 2 }
func (f *statFile) Fulfill(ctx context.Context, dispatchErr func(error)) error {
	r, w, err := os.Pipe()
	if err != nil {
		return err
	}
	f.Set(w)

	done := make(chan struct{})
	go func() {
		defer close(done)
		var n int

		n, err = r.Read(make([]byte, 1))
		switch n {
		case -1:
			if err == nil {
				err = ErrStatFault
			}
			dispatchErr(err)
		case 0:
			if err == nil {
				err = ErrStatRead
			}
			dispatchErr(err)
		case 1:
			dispatchErr(err)
		default:
			panic("unreachable")
		}
		runtime.KeepAlive(w)
	}()

	go func() {
		select {
		case <-done:
			dispatchErr(nil)
		case <-ctx.Done():
			dispatchErr(r.Close()) // this aborts Read with file already closed
			runtime.KeepAlive(w)
		}
	}()

	// this gets closed by the caller
	*f.s = r
	return nil
}