-
Notifications
You must be signed in to change notification settings - Fork 3.8k
/
KubePodProcessIntegrationTest.java
207 lines (177 loc) · 6.71 KB
/
KubePodProcessIntegrationTest.java
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
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
/*
* MIT License
*
* Copyright (c) 2020 Airbyte
*
* Permission is hereby granted, free of charge, to any person obtaining a copy
* of this software and associated documentation files (the "Software"), to deal
* in the Software without restriction, including without limitation the rights
* to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
* copies of the Software, and to permit persons to whom the Software is
* furnished to do so, subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
* FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
* AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
* LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
* OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
* SOFTWARE.
*/
package io.airbyte.workers.process;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import com.google.common.collect.ImmutableMap;
import io.airbyte.commons.lang.Exceptions;
import io.airbyte.workers.WorkerException;
import io.fabric8.kubernetes.client.DefaultKubernetesClient;
import io.fabric8.kubernetes.client.KubernetesClient;
import io.kubernetes.client.openapi.ApiClient;
import io.kubernetes.client.util.Config;
import java.io.IOException;
import java.net.Inet4Address;
import java.net.ServerSocket;
import java.net.UnknownHostException;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingDeque;
import org.apache.commons.lang.RandomStringUtils;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
// requires kube running locally to run. If using Minikube it requires MINIKUBE=true
public class KubePodProcessIntegrationTest {
private static final boolean IS_MINIKUBE = Boolean.parseBoolean(Optional.ofNullable(System.getenv("IS_MINIKUBE")).orElse("false"));
private List<Integer> openPorts;
private List<Integer> openWorkerPorts;
private int heartbeatPort;
private String heartbeatUrl;
private ApiClient officialClient;
private KubernetesClient fabricClient;
private BlockingQueue<Integer> workerPorts;
private KubeProcessFactory processFactory;
private static WorkerHeartbeatServer server;
@BeforeEach
public void setup() throws Exception {
openPorts = new ArrayList<>(getOpenPorts(5));
openWorkerPorts = openPorts.subList(1, openPorts.size() - 1);
heartbeatPort = openPorts.get(0);
heartbeatUrl = getHost() + ":" + heartbeatPort;
officialClient = Config.defaultClient();
fabricClient = new DefaultKubernetesClient();
workerPorts = new LinkedBlockingDeque<>(openWorkerPorts);
processFactory = new KubeProcessFactory("default", officialClient, fabricClient, heartbeatUrl, workerPorts);
server = new WorkerHeartbeatServer(heartbeatPort);
server.startBackground();
}
@AfterEach
public void teardown() throws Exception {
server.stop();
}
@Test
public void testSuccessfulSpawning() throws Exception {
// start a finite process
final Process process = getProcess("echo hi; sleep 1; echo hi2");
process.waitFor();
// the pod should be dead and in a good state
assertFalse(process.isAlive());
assertEquals(0, process.exitValue());
}
@Test
public void testPipeInEntrypoint() throws Exception {
// start a process that has a pipe in the entrypoint
final Process process = getProcess("echo hi | cat");
process.waitFor();
// the pod should be dead and in a good state
assertFalse(process.isAlive());
assertEquals(0, process.exitValue());
}
@Test
public void testExitCodeRetrieval() throws Exception {
// start a process that requests
final Process process = getProcess("exit 10");
process.waitFor();
// the pod should be dead with the correct error code
assertFalse(process.isAlive());
assertEquals(10, process.exitValue());
}
@Test
public void testMissingEntrypoint() throws WorkerException, InterruptedException {
// start a process with an entrypoint that doesn't exist
final Process process = getProcess("ksaiiiasdfjklaslkei");
process.waitFor();
// the pod should be dead and in an error state
assertFalse(process.isAlive());
assertEquals(127, process.exitValue());
}
@Test
public void testKillingWithoutHeartbeat() throws Exception {
// start an infinite process
final Process process = getProcess("while true; do echo hi; sleep 1; done");
// kill the heartbeat server
server.stop();
// waiting for process
process.waitFor();
// the pod should be dead and in an error state
assertFalse(process.isAlive());
assertNotEquals(0, process.exitValue());
}
private static String getRandomFile(int lines) {
var sb = new StringBuilder();
for (int i = 0; i < lines; i++) {
sb.append(RandomStringUtils.randomAlphabetic(100));
sb.append("\n");
}
return sb.toString();
}
private Process getProcess(String entrypoint) throws WorkerException {
// these files aren't used for anything, it's just to check for exceptions when uploading
var files = ImmutableMap.of(
"file0", "fixed str",
"file1", getRandomFile(1),
"file2", getRandomFile(100),
"file3", getRandomFile(1000));
return processFactory.create(
"some-id",
0,
Path.of("/tmp/job-root"),
"busybox:latest",
false,
files,
entrypoint);
}
private static Set<Integer> getOpenPorts(int count) {
final Set<ServerSocket> servers = new HashSet<>();
final Set<Integer> ports = new HashSet<>();
try {
for (int i = 0; i < count; i++) {
var server = new ServerSocket(0);
servers.add(server);
ports.add(server.getLocalPort());
}
} catch (IOException e) {
throw new RuntimeException(e);
} finally {
for (ServerSocket server : servers) {
Exceptions.swallow(server::close);
}
}
return ports;
}
private static String getHost() {
try {
return (IS_MINIKUBE ? Inet4Address.getLocalHost().getHostAddress() : "host.docker.internal");
} catch (UnknownHostException e) {
throw new RuntimeException(e);
}
};
}