-
Notifications
You must be signed in to change notification settings - Fork 34
Expand file tree
/
Copy pathservice_lifecycle_test.go
More file actions
144 lines (128 loc) · 4 KB
/
Copy pathservice_lifecycle_test.go
File metadata and controls
144 lines (128 loc) · 4 KB
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
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
package kata
import (
"bufio"
"context"
"io"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.kenn.io/kata/internal/daemon"
"go.kenn.io/kata/internal/db"
)
func TestServiceCloseTerminatesActiveEventStream(t *testing.T) {
service, err := New(context.Background(), Config{
DSN: filepath.Join(t.TempDir(), "service.db"),
Auth: AuthConfig{TrustCallerAuthentication: true},
})
require.NoError(t, err)
server := httptest.NewServer(service.Handler())
t.Cleanup(server.Close)
request, err := http.NewRequest(http.MethodGet, server.URL+"/api/v1/events/stream?after_id=0", nil)
require.NoError(t, err)
request.Header.Set("Accept", "text/event-stream")
response, err := server.Client().Do(request)
require.NoError(t, err)
defer func() { _ = response.Body.Close() }()
require.Equal(t, http.StatusOK, response.StatusCode)
project, err := service.store.CreateProject(context.Background(), "example-project")
require.NoError(t, err)
_, event, err := service.store.CreateIssue(context.Background(), db.CreateIssueParams{
ProjectID: project.ID,
Title: "observe shutdown",
Author: "example-user",
})
require.NoError(t, err)
service.broadcaster.Broadcast(daemon.StreamMsg{
Kind: "event", Event: &event, ProjectID: project.ID,
})
reader := bufio.NewReader(response.Body)
eventSeen := make(chan error, 1)
go func() {
for {
line, readErr := reader.ReadString('\n')
if readErr != nil {
eventSeen <- readErr
return
}
if strings.TrimSpace(line) == "event: issue.created" {
eventSeen <- nil
return
}
}
}()
select {
case eventErr := <-eventSeen:
require.NoError(t, eventErr)
case <-time.After(2 * time.Second):
require.Fail(t, "event stream did not enter its live phase")
}
streamDone := make(chan error, 1)
go func() {
_, copyErr := io.Copy(io.Discard, reader)
streamDone <- copyErr
}()
closeDone := make(chan error, 1)
go func() { closeDone <- service.Close() }()
select {
case closeErr := <-closeDone:
require.NoError(t, closeErr)
case <-time.After(2 * time.Second):
require.Fail(t, "Close did not wait for the active event stream")
}
select {
case streamErr := <-streamDone:
require.NoError(t, streamErr)
case <-time.After(2 * time.Second):
require.Fail(t, "active event stream outlived Close")
}
postClose, err := server.Client().Get(server.URL + "/api/v1/health")
require.NoError(t, err)
defer func() { _ = postClose.Body.Close() }()
assert.Equal(t, http.StatusServiceUnavailable, postClose.StatusCode)
}
func TestServiceEnsureProjectBroadcastsAndEnqueuesExactCreatedEvent(t *testing.T) {
service, err := New(context.Background(), Config{
DSN: filepath.Join(t.TempDir(), "service.db"),
Auth: AuthConfig{TrustCallerAuthentication: true},
})
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, service.Close()) })
sink := &serviceRecordingSink{}
service.hookSink = sink
subscription := service.broadcaster.Subscribe(daemon.SubFilter{})
defer subscription.Unsub()
result, err := service.EnsureProject(context.Background(), ProjectSpec{
UID: "01HZNQ7VFPK1XGD8R5MABCD4EX", Name: "example-project",
})
require.NoError(t, err)
require.True(t, result.Created)
var msg daemon.StreamMsg
select {
case msg = <-subscription.Ch:
case <-time.After(time.Second):
t.Fatal("project creation event was not broadcast")
}
require.Equal(t, "event", msg.Kind)
require.NotNil(t, msg.Event)
assert.Equal(t, "project.created", msg.Event.Type)
assert.Equal(t, db.SystemActor, msg.Event.Actor)
require.Len(t, sink.events, 1)
assert.Equal(t, *msg.Event, sink.events[0])
persisted, err := service.store.EventsAfter(context.Background(), db.EventsAfterParams{
ProjectID: result.Project.ID, Limit: 10,
})
require.NoError(t, err)
require.Len(t, persisted, 1)
assert.Equal(t, *msg.Event, persisted[0])
}
type serviceRecordingSink struct {
events []db.Event
}
func (s *serviceRecordingSink) Enqueue(event db.Event) {
s.events = append(s.events, event)
}