Skip to content

Commit 2483cba

Browse files
committed
Continue unifying the new session queue
We have now made the `SessionRequests` package-private, so it's only visible to the `LocalNewSessionQueue`. At this point, we've isolated the rest of the Grid from the implementation of the session queue.
1 parent 1d31428 commit 2483cba

19 files changed

Lines changed: 248 additions & 310 deletions

java/server/src/org/openqa/selenium/grid/commands/Hub.java

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
import com.google.auto.service.AutoService;
2121
import com.google.common.collect.ImmutableSet;
22-
2322
import org.openqa.selenium.BuildInfo;
2423
import org.openqa.selenium.cli.CliCommand;
2524
import org.openqa.selenium.events.EventBus;
@@ -41,10 +40,8 @@
4140
import org.openqa.selenium.grid.server.Server;
4241
import org.openqa.selenium.grid.sessionmap.SessionMap;
4342
import org.openqa.selenium.grid.sessionmap.local.LocalSessionMap;
44-
import org.openqa.selenium.grid.sessionqueue.local.SessionRequests;
4543
import org.openqa.selenium.grid.sessionqueue.NewSessionQueue;
4644
import org.openqa.selenium.grid.sessionqueue.config.SessionRequestOptions;
47-
import org.openqa.selenium.grid.sessionqueue.local.SessionRequests;
4845
import org.openqa.selenium.grid.sessionqueue.local.LocalNewSessionQueue;
4946
import org.openqa.selenium.grid.web.CombinedHandler;
5047
import org.openqa.selenium.grid.web.GridUiRoute;
@@ -144,12 +141,12 @@ protected Handlers createHandlers(Config config) {
144141
networkOptions.getHttpClientFactory(tracer));
145142

146143
SessionRequestOptions sessionRequestOptions = new SessionRequestOptions(config);
147-
SessionRequests sessionRequests = new SessionRequests(
144+
NewSessionQueue queue = new LocalNewSessionQueue(
148145
tracer,
149-
bus,
146+
bus,
150147
sessionRequestOptions.getSessionRequestRetryInterval(),
151-
sessionRequestOptions.getSessionRequestTimeout());
152-
NewSessionQueue queue = new LocalNewSessionQueue(tracer, bus, sessionRequests, secret);
148+
sessionRequestOptions.getSessionRequestTimeout(),
149+
secret);
153150
handler.addHandler(queue);
154151

155152
DistributorOptions distributorOptions = new DistributorOptions(config);

java/server/src/org/openqa/selenium/grid/commands/Standalone.java

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,6 @@
4242
import org.openqa.selenium.grid.server.Server;
4343
import org.openqa.selenium.grid.sessionmap.SessionMap;
4444
import org.openqa.selenium.grid.sessionmap.local.LocalSessionMap;
45-
import org.openqa.selenium.grid.sessionqueue.local.SessionRequests;
4645
import org.openqa.selenium.grid.sessionqueue.NewSessionQueue;
4746
import org.openqa.selenium.grid.sessionqueue.config.SessionRequestOptions;
4847
import org.openqa.selenium.grid.sessionqueue.local.LocalNewSessionQueue;
@@ -140,16 +139,11 @@ protected Handlers createHandlers(Config config) {
140139
combinedHandler.addHandler(sessions);
141140

142141
SessionRequestOptions sessionRequestOptions = new SessionRequestOptions(config);
143-
SessionRequests sessionRequests = new SessionRequests(
144-
tracer,
145-
bus,
146-
sessionRequestOptions.getSessionRequestRetryInterval(),
147-
sessionRequestOptions.getSessionRequestTimeout());
148-
149142
NewSessionQueue queue = new LocalNewSessionQueue(
150143
tracer,
151144
bus,
152-
sessionRequests,
145+
sessionRequestOptions.getSessionRequestRetryInterval(),
146+
sessionRequestOptions.getSessionRequestTimeout(),
153147
registrationSecret);
154148
combinedHandler.addHandler(queue);
155149

java/server/src/org/openqa/selenium/grid/sessionqueue/NewSessionQueue.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,8 @@ protected NewSessionQueue(Tracer tracer, Secret registrationSecret) {
7373
throw new UncheckedIOException(e);
7474
}
7575
}),
76+
post("/se/grid/newsessionqueue/session/last")
77+
.to(() -> new OfferLastToSessionQueue(tracer, this)),
7678
post("/se/grid/newsessionqueue/session")
7779
.to(() -> new AddToSessionQueue(tracer, this)),
7880
post("/se/grid/newsessionqueue/session/retry/{requestId}")
@@ -94,6 +96,8 @@ private RequestId requestIdFrom(Map<String, String> params) {
9496

9597
public abstract HttpResponse addToQueue(SessionRequest request);
9698

99+
public abstract boolean offerLast(SessionRequest request);
100+
97101
public abstract boolean retryAddToQueue(SessionRequest request);
98102

99103
public abstract Optional<SessionRequest> remove(RequestId reqId);
@@ -111,6 +115,5 @@ public boolean matches(HttpRequest req) {
111115
public HttpResponse execute(HttpRequest req) {
112116
return routes.execute(req);
113117
}
114-
115118
}
116119

Lines changed: 61 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,61 @@
1+
// Licensed to the Software Freedom Conservancy (SFC) under one
2+
// or more contributor license agreements. See the NOTICE file
3+
// distributed with this work for additional information
4+
// regarding copyright ownership. The SFC licenses this file
5+
// to you under the Apache License, Version 2.0 (the
6+
// "License"); you may not use this file except in compliance
7+
// with the License. You may obtain a copy of the License at
8+
//
9+
// http://www.apache.org/licenses/LICENSE-2.0
10+
//
11+
// Unless required by applicable law or agreed to in writing,
12+
// software distributed under the License is distributed on an
13+
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
// KIND, either express or implied. See the License for the
15+
// specific language governing permissions and limitations
16+
// under the License.
17+
18+
package org.openqa.selenium.grid.sessionqueue;
19+
20+
import org.openqa.selenium.internal.Require;
21+
import org.openqa.selenium.remote.http.Contents;
22+
import org.openqa.selenium.remote.http.HttpHandler;
23+
import org.openqa.selenium.remote.http.HttpRequest;
24+
import org.openqa.selenium.remote.http.HttpResponse;
25+
import org.openqa.selenium.remote.tracing.Span;
26+
import org.openqa.selenium.remote.tracing.Tracer;
27+
28+
import java.util.Collections;
29+
30+
import static org.openqa.selenium.remote.http.Contents.asJson;
31+
import static org.openqa.selenium.remote.tracing.HttpTracing.newSpanAsChildOf;
32+
import static org.openqa.selenium.remote.tracing.Tags.HTTP_REQUEST;
33+
import static org.openqa.selenium.remote.tracing.Tags.HTTP_RESPONSE;
34+
35+
class OfferLastToSessionQueue implements HttpHandler {
36+
37+
private final Tracer tracer;
38+
private final NewSessionQueue newSessionQueue;
39+
40+
OfferLastToSessionQueue(Tracer tracer, NewSessionQueue newSessionQueue) {
41+
this.tracer = Require.nonNull("Tracer", tracer);
42+
this.newSessionQueue = Require.nonNull("New Session Queue", newSessionQueue);
43+
}
44+
45+
@Override
46+
public HttpResponse execute(HttpRequest req) {
47+
try (Span span = newSpanAsChildOf(tracer, req, "sessionqueue.addLast")) {
48+
HTTP_REQUEST.accept(span, req);
49+
50+
boolean result = newSessionQueue.offerLast(Contents.fromJson(req, SessionRequest.class));
51+
52+
HttpResponse response = new HttpResponse()
53+
.setContent(Contents.asJson(Collections.singletonMap("value", result)));
54+
55+
HTTP_RESPONSE.accept(span, response);
56+
57+
return response;
58+
}
59+
}
60+
}
61+

java/server/src/org/openqa/selenium/grid/sessionqueue/local/LocalNewSessionQueue.java

Lines changed: 27 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,9 @@
2222
import org.openqa.selenium.events.EventBus;
2323
import org.openqa.selenium.grid.config.Config;
2424
import org.openqa.selenium.grid.data.RequestId;
25+
import org.openqa.selenium.grid.jmx.JMXHelper;
26+
import org.openqa.selenium.grid.jmx.ManagedAttribute;
27+
import org.openqa.selenium.grid.jmx.ManagedService;
2528
import org.openqa.selenium.grid.log.LoggingOptions;
2629
import org.openqa.selenium.grid.security.Secret;
2730
import org.openqa.selenium.grid.security.SecretOptions;
@@ -44,9 +47,10 @@
4447
import java.util.Optional;
4548
import java.util.Set;
4649

47-
import static org.openqa.selenium.remote.http.Contents.reader;
4850
import static org.openqa.selenium.remote.tracing.Tags.EXCEPTION;
4951

52+
@ManagedService(objectName = "org.seleniumhq.grid:type=SessionQueue,name=LocalSessionQueue",
53+
description = "New session queue")
5054
public class LocalNewSessionQueue extends NewSessionQueue {
5155

5256
public final SessionRequests sessionRequests;
@@ -56,30 +60,33 @@ public class LocalNewSessionQueue extends NewSessionQueue {
5660
public LocalNewSessionQueue(
5761
Tracer tracer,
5862
EventBus bus,
59-
SessionRequests sessionRequests,
63+
Duration retryInterval,
64+
Duration requestTimeout,
6065
Secret registrationSecret) {
6166
super(tracer, registrationSecret);
6267
this.bus = Require.nonNull("Event bus", bus);
63-
this.sessionRequests = Require.nonNull("New Session Request Queue", sessionRequests);
68+
69+
this.sessionRequests = new SessionRequests(
70+
tracer,
71+
bus,
72+
Require.nonNull("Retry interval", retryInterval),
73+
Require.nonNull("Request timeout", requestTimeout));
6474

6575
this.getNewSessionResponse = new GetNewSessionResponse(bus, sessionRequests);
76+
77+
new JMXHelper().register(this);
6678
}
6779

6880
public static NewSessionQueue create(Config config) {
6981
Tracer tracer = new LoggingOptions(config).getTracer();
7082
EventBus bus = new EventBusOptions(config).getEventBus();
7183
Duration retryInterval = new SessionRequestOptions(config).getSessionRequestRetryInterval();
7284
Duration requestTimeout = new SessionRequestOptions(config).getSessionRequestTimeout();
73-
SessionRequests sessionRequests = new SessionRequests(
74-
tracer,
75-
bus,
76-
retryInterval,
77-
requestTimeout);
7885

7986
SecretOptions secretOptions = new SecretOptions(config);
8087
Secret registrationSecret = secretOptions.getRegistrationSecret();
8188

82-
return new LocalNewSessionQueue(tracer, bus, sessionRequests, registrationSecret);
89+
return new LocalNewSessionQueue(tracer, bus, retryInterval, requestTimeout, registrationSecret);
8390
}
8491

8592
@Override
@@ -88,6 +95,12 @@ public HttpResponse addToQueue(SessionRequest request) {
8895
return getNewSessionResponse.add(request);
8996
}
9097

98+
@Override
99+
public boolean offerLast(SessionRequest request) {
100+
Require.nonNull("Session request", request);
101+
return sessionRequests.offerLast(request);
102+
}
103+
91104
@Override
92105
public boolean retryAddToQueue(SessionRequest request) {
93106
return sessionRequests.offerFirst(request);
@@ -113,6 +126,11 @@ public boolean isReady() {
113126
return bus.isReady();
114127
}
115128

129+
@ManagedAttribute(name = "NewSessionQueueSize")
130+
public int getQueueSize() {
131+
return sessionRequests.getQueueSize();
132+
}
133+
116134
private void validateSessionRequest(SessionRequest request) {
117135
try (Span span = tracer.getCurrentContext().createSpan("newsession_queue.validate")) {
118136
Map<String, EventAttributeValue> attributeMap = new HashMap<>();

java/server/src/org/openqa/selenium/grid/sessionqueue/local/SessionRequests.java

Lines changed: 14 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,6 @@
2525
import org.openqa.selenium.grid.data.NewSessionRejectedEvent;
2626
import org.openqa.selenium.grid.data.NewSessionRequestEvent;
2727
import org.openqa.selenium.grid.data.RequestId;
28-
import org.openqa.selenium.grid.jmx.JMXHelper;
29-
import org.openqa.selenium.grid.jmx.ManagedAttribute;
30-
import org.openqa.selenium.grid.jmx.ManagedService;
3128
import org.openqa.selenium.grid.log.LoggingOptions;
3229
import org.openqa.selenium.grid.server.EventBusOptions;
3330
import org.openqa.selenium.grid.sessionqueue.SessionRequest;
@@ -59,38 +56,39 @@
5956
import java.util.logging.Logger;
6057
import java.util.stream.Collectors;
6158

62-
@ManagedService(objectName = "org.seleniumhq.grid:type=SessionQueue,name=LocalSessionQueue",
63-
description = "New session queue")
64-
public class SessionRequests {
59+
class SessionRequests {
6560

6661
private static final Logger LOG = Logger.getLogger(SessionRequests.class.getName());
67-
private final EventBus bus;
6862
private final Tracer tracer;
63+
private final EventBus bus;
6964
private final Duration retryInterval;
7065
private final Duration requestTimeout;
7166
private final Deque<SessionRequest> sessionRequests = new ConcurrentLinkedDeque<>();
7267
private final ReadWriteLock lock = new ReentrantReadWriteLock(true);
7368
private final ScheduledExecutorService executorService =
7469
Executors.newSingleThreadScheduledExecutor();
7570
private final Thread shutdownHook = new Thread(this::callExecutorShutdown);
76-
71+
private final String timedOutErrorMessage;
7772

7873
public SessionRequests(
7974
Tracer tracer,
8075
EventBus bus,
8176
Duration retryInterval,
8277
Duration requestTimeout) {
8378
this.tracer = Require.nonNull("Tracer", tracer);
79+
this.bus = Require.nonNull("Event bus", bus);
8480
this.retryInterval = Require.nonNull("Session request retry interval", retryInterval);
8581
this.requestTimeout = Require.nonNull("Session request timeout", requestTimeout);
86-
this.bus = Require.nonNull("Event bus", bus);
82+
83+
timedOutErrorMessage = String.format(
84+
"New session request rejected after being in the queue for more than %s",
85+
format(requestTimeout));
86+
8787
Runtime.getRuntime().addShutdownHook(shutdownHook);
8888

8989
Regularly regularly = new Regularly("New Session Queue Clean up");
9090
Duration purgeTimeout = requestTimeout.multipliedBy(2);
9191
regularly.submit(this::purgeTimedOutRequests, purgeTimeout, purgeTimeout);
92-
93-
new JMXHelper().register(this);
9492
}
9593

9694
public static SessionRequests create(Config config) {
@@ -105,8 +103,7 @@ public boolean isReady() {
105103
return bus.isReady();
106104
}
107105

108-
@ManagedAttribute(name = "NewSessionQueueSize")
109-
public int getQueueSize() {
106+
int getQueueSize() {
110107
Lock readLock = lock.readLock();
111108
readLock.lock();
112109
try {
@@ -166,8 +163,7 @@ public boolean offerFirst(SessionRequest request) {
166163
try {
167164
boolean added = sessionRequests.offerFirst(request);
168165
if (added) {
169-
executorService.schedule(() -> retryRequest(request),
170-
retryInterval.getSeconds(), TimeUnit.SECONDS);
166+
executorService.schedule(() -> retryRequest(request), retryInterval.getSeconds(), TimeUnit.SECONDS);
171167
}
172168
return added;
173169
} finally {
@@ -184,7 +180,7 @@ private void retryRequest(SessionRequest sessionRequest) {
184180
LOG.log(Level.INFO, "Request {0} timed out", requestId);
185181
sessionRequests.remove(sessionRequest);
186182
bus.fire(new NewSessionRejectedEvent(
187-
new NewSessionErrorResponse(requestId, getTimeoutErrorMessage())));
183+
new NewSessionErrorResponse(requestId, timedOutErrorMessage)));
188184
} else {
189185
LOG.log(Level.INFO,
190186
"Adding request back to the queue. All slots are busy. Request: {0}",
@@ -226,7 +222,7 @@ public Optional<SessionRequest> remove(RequestId id) {
226222
if (request.isPresent()) {
227223
if (hasRequestTimedOut(request.get())) {
228224
bus.fire(new NewSessionRejectedEvent(
229-
new NewSessionErrorResponse(id, getTimeoutErrorMessage())));
225+
new NewSessionErrorResponse(id, timedOutErrorMessage)));
230226
return Optional.empty();
231227
}
232228
}
@@ -267,7 +263,7 @@ private void purgeTimedOutRequests() {
267263
if (hasRequestTimedOut(sessionRequest)) {
268264
iterator.remove();
269265
bus.fire(new NewSessionRejectedEvent(
270-
new NewSessionErrorResponse(sessionRequest.getRequestId(), getTimeoutErrorMessage())));
266+
new NewSessionErrorResponse(sessionRequest.getRequestId(), timedOutErrorMessage)));
271267
}
272268
}
273269
} finally {
@@ -303,10 +299,4 @@ private static String format(Duration duration) {
303299
toReturn.append(secs).append("s");
304300
return toReturn.toString();
305301
}
306-
307-
private String getTimeoutErrorMessage() {
308-
return String.format(
309-
"New session request rejected after being in the queue for more than %s",
310-
format(requestTimeout));
311-
}
312302
}

java/server/src/org/openqa/selenium/grid/sessionqueue/remote/BUILD.bazel

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ java_library(
77
visibility = [
88
"//java/server/src/org/openqa/selenium/grid:__subpackages__",
99
"//java/server/test/org/openqa/selenium/grid/router:__pkg__",
10-
"//java/server/test/org/openqa/selenium/grid/sessionqueue:__pkg__",
10+
"//java/server/test/org/openqa/selenium/grid/sessionqueue:__subpackages__",
1111
],
1212
deps = [
1313
"//java/client/src/org/openqa/selenium/json",

java/server/src/org/openqa/selenium/grid/sessionqueue/remote/RemoteNewSessionQueue.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,15 @@ public HttpResponse addToQueue(SessionRequest request) {
9393
return client.execute(upstream);
9494
}
9595

96+
@Override
97+
public boolean offerLast(SessionRequest request) {
98+
HttpRequest upstream = new HttpRequest(POST, "/se/grid/newsessionqueue/session/last");
99+
HttpTracing.inject(tracer, tracer.getCurrentContext(), upstream);
100+
upstream.setContent(Contents.asJson(request));
101+
HttpResponse response = client.execute(upstream);
102+
return Values.get(response, Boolean.class);
103+
}
104+
96105
@Override
97106
public boolean retryAddToQueue(SessionRequest request) {
98107
Require.nonNull("Session request", request);

0 commit comments

Comments
 (0)