Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@

package com.volcengine.ark.runtime.selfhosted;

import com.volcengine.ark.runtime.models.environment.HeartbeatWorkResponse;
import com.volcengine.ark.runtime.models.environment.WorkItem;
import com.volcengine.ark.runtime.models.environment.WorkState;
import java.io.IOException;
import java.lang.management.ManagementFactory;
import java.net.InetAddress;
Expand Down Expand Up @@ -58,7 +61,7 @@ public void run() {
return;
}
try {
handleItem(item, false);
handleItem(claimedWorkFromItem(item), false);
} catch (SessionToolRunner.IdleTimeoutException | SessionToolRunner.SessionTerminatedException ignored) {
} catch (Exception e) {
options.logger.log(Level.WARNING, "handle work failed", e);
Expand All @@ -74,52 +77,45 @@ public void handleItem(HandleItemOptions handleOptions) throws IOException {
Thread previous = activeThread;
activeThread = Thread.currentThread();
try {
handleItem(workItemFromOptions(handleOptions), true);
handleItem(claimedWorkFromOptions(handleOptions), true);
} catch (SessionToolRunner.IdleTimeoutException | SessionToolRunner.SessionTerminatedException ignored) {
} finally {
activeThread = previous;
}
}

private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOException {
if (item.getEnvironmentId() == null || item.getEnvironmentId().isEmpty()) {
item.setEnvironmentId(firstNonEmpty(options.environmentId, System.getenv("MA_ENVIRONMENT_ID")));
}
if (item.getId() == null || item.getId().isEmpty()) {
throw new IllegalArgumentException("work item id must not be empty");
}
String sessionId = item.sessionIdValue();
if (sessionId.isEmpty()) {
throw new IllegalArgumentException("work item does not contain session id");
private void handleItem(ClaimedWork work, boolean useWorkdirAsSession) throws IOException {
if (work.environmentId.isEmpty()) {
work.environmentId = firstNonEmpty(options.environmentId, System.getenv("MA_ENVIRONMENT_ID"));
}
AtomicBoolean stop = new AtomicBoolean(false);
AtomicReference<String> heartbeatCause = new AtomicReference<>("");
Thread heartbeat = null;
try {
String workdir = workdirFor(sessionId, useWorkdirAsSession);
String workdir = workdirFor(work.sessionId, useWorkdirAsSession);
Thread heartbeatThread = new Thread(
() -> heartbeatLoop(item, stop, heartbeatCause), "ma-self-host-heartbeat");
() -> heartbeatLoop(work, stop, heartbeatCause), "ma-self-host-heartbeat");
heartbeatThread.setDaemon(true);
heartbeatThread.start();
heartbeat = heartbeatThread;
SessionSnapshot session = api.getSession(sessionId);
SessionSnapshot session = api.getSession(work.sessionId);
if (closed.get() || stop.get()) {
return;
}
if (session == null) {
throw new IOException("session response is empty");
}
if (session.getId() == null || session.getId().isEmpty()) {
session.setId(sessionId);
session.setId(work.sessionId);
}
new Initializer(api, new Initializer.Options(workdir)).setup(session);
if (closed.get() || stop.get()) {
return;
}
ToolContext toolContext = toolContext(workdir, stop);
FileToolResultStore store = new FileToolResultStore(workdir);
SessionToolRunner runner = new SessionToolRunner(api, sessionId, new SessionToolRunner.Options()
.workId(item.getId())
SessionToolRunner runner = new SessionToolRunner(api, work.sessionId, new SessionToolRunner.Options()
.workId(work.id)
.tools(options.tools == null ? DefaultTools.create() : options.tools)
.toolContext(toolContext)
.customTools(options.customTools)
Expand All @@ -145,7 +141,7 @@ private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOExc
String cause = heartbeatCause.get();
if (shouldStopItem(cause)) {
try {
api.stopWork(item.getEnvironmentId(), item.getId(), true);
api.stopWork(work.environmentId, work.id, true);
} catch (RuntimeException e) {
if (!isResolvedStatus(e)) {
options.logger.log(Level.WARNING, "stop work failed", e);
Expand All @@ -158,19 +154,19 @@ private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOExc
}
}

private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference<String> cause) {
private void heartbeatLoop(ClaimedWork work, AtomicBoolean stop, AtomicReference<String> cause) {
long interval = Math.max(1000L, SelfHostedConstants.DEFAULT_HEARTBEAT_MILLIS / 2L);
long ttl = SelfHostedConstants.DEFAULT_HEARTBEAT_MILLIS;
String last = item.latestHeartbeatValue();
String last = work.latestHeartbeatAt;
if (last == null || last.isEmpty()) {
last = SelfHostedConstants.EXPECTED_LAST_HEARTBEAT_NO_HEARTBEAT;
}
long lastSuccess = System.currentTimeMillis();
while (!stop.get()) {
try {
HeartbeatResponse response = api.heartbeatWork(
item.getEnvironmentId(),
item.getId(),
HeartbeatWorkResponse response = api.heartbeatWork(
work.environmentId,
work.id,
last,
(int) (ttl / 1000L));
if (response == null) {
Expand All @@ -180,21 +176,13 @@ private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference<St
return;
}
options.logger.warning(
"heartbeat empty response work_id=" + item.getId()
+ " session_id=" + item.sessionIdValue());
"heartbeat empty response work_id=" + work.id
+ " session_id=" + work.sessionId);
sleep(interval, stop);
continue;
}
lastSuccess = System.currentTimeMillis();
if (response.getLastHeartbeat() != null && !response.getLastHeartbeat().isEmpty()) {
last = response.getLastHeartbeat();
}
if (response.getTtlSeconds() > 0) {
ttl = response.getTtlSeconds() * 1000L;
interval = Math.max(1000L, Math.min(ttl / 2, SelfHostedConstants.DEFAULT_HEARTBEAT_MILLIS));
}
if (SelfHostedConstants.WORK_STATE_STOPPING.equals(response.getState())
|| SelfHostedConstants.WORK_STATE_STOPPED.equals(response.getState())) {
if (WorkState.STOPPING.equals(response.getState())
|| WorkState.STOPPED.equals(response.getState())) {
cause.set("stop_requested");
stop.set(true);
return;
Expand All @@ -204,6 +192,15 @@ private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference<St
stop.set(true);
return;
}
lastSuccess = System.currentTimeMillis();
if (response.getLastHeartbeat() != null && !response.getLastHeartbeat().isEmpty()) {
last = response.getLastHeartbeat();
}
Long ttlSeconds = response.getTtlSeconds();
if (ttlSeconds != null && ttlSeconds > 0) {
ttl = ttlSeconds * 1000L;
interval = Math.max(1000L, Math.min(ttl / 2, SelfHostedConstants.DEFAULT_HEARTBEAT_MILLIS));
}
} catch (RuntimeException e) {
if (WorkerAPIException.isStatus(e, 412)) {
cause.set("lease_lost");
Expand All @@ -222,8 +219,8 @@ private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference<St
}
options.logger.log(
Level.WARNING,
"heartbeat failed work_id=" + item.getId()
+ " session_id=" + item.sessionIdValue()
"heartbeat failed work_id=" + work.id
+ " session_id=" + work.sessionId
+ " since_last_success_ms=" + (System.currentTimeMillis() - lastSuccess)
+ " ttl_ms=" + ttl,
e);
Expand Down Expand Up @@ -257,7 +254,7 @@ private String workdirFor(String sessionId, boolean useWorkdirAsSession) throws
return sessionDir.toString();
}

private WorkItem workItemFromOptions(HandleItemOptions opts) {
private ClaimedWork claimedWorkFromOptions(HandleItemOptions opts) {
String workId = firstNonEmpty(opts.workId, System.getenv("MA_WORK_ID"));
String environmentId = firstNonEmpty(opts.environmentId, System.getenv("MA_ENVIRONMENT_ID"));
String sessionId = firstNonEmpty(opts.sessionId, System.getenv("MA_SESSION_ID"));
Expand All @@ -271,15 +268,22 @@ private WorkItem workItemFromOptions(HandleItemOptions opts) {
if (sessionId.isEmpty()) {
throw new IllegalArgumentException("session id is required");
}
WorkData data = new WorkData();
data.setType("session");
data.setId(sessionId);
WorkItem item = new WorkItem();
item.setId(workId);
item.setEnvironmentId(environmentId);
item.setLatestHeartbeatAt(latestHeartbeat);
item.setData(data);
return item;
return new ClaimedWork(workId, environmentId, sessionId, latestHeartbeat);
}

private static ClaimedWork claimedWorkFromItem(WorkItem item) {
if (item == null || item.getId() == null || item.getId().isEmpty()) {
throw new IllegalArgumentException("work item id must not be empty");
}
String sessionId = WorkItems.sessionId(item);
if (sessionId.isEmpty()) {
throw new IllegalArgumentException("work item does not contain session id");
}
return new ClaimedWork(
item.getId(),
item.getEnvironmentId() == null ? "" : item.getEnvironmentId(),
sessionId,
WorkItems.latestHeartbeat(item));
}

private static void sleep(long millis, AtomicBoolean stop) {
Expand Down Expand Up @@ -357,6 +361,20 @@ private static String firstNonEmpty(String first, String second) {
return first != null && !first.isEmpty() ? first : (second == null ? "" : second);
}

private static final class ClaimedWork {
private final String id;
private String environmentId;
private final String sessionId;
private final String latestHeartbeatAt;

private ClaimedWork(String id, String environmentId, String sessionId, String latestHeartbeatAt) {
this.id = id;
this.environmentId = environmentId;
this.sessionId = sessionId;
this.latestHeartbeatAt = latestHeartbeatAt;
}
}

public static class HandleItemOptions {
private String workId = "";
private String environmentId = "";
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import com.volcengine.ark.runtime.models.environment.EnvironmentWorkPoll200Response;
import com.volcengine.ark.runtime.models.environment.HeartbeatWorkResponse;
import com.volcengine.ark.runtime.models.environment.StopWorkBody;
import com.volcengine.ark.runtime.models.environment.WorkItem;
import com.volcengine.ark.runtime.models.session.ManagedAgentsEventParams;
import com.volcengine.ark.runtime.models.session.SendSessionEventsRequest;
import com.volcengine.ark.runtime.models.skill.Skill;
Expand Down Expand Up @@ -100,7 +101,7 @@ public WorkItem pollWork(String environmentId, String workerId, int blockMs, int
if (response == null || response.getId() == null || response.getId().isEmpty()) {
return null;
}
return WorkItem.fromMap(toMap(response));
return mapper.convertValue(response, WorkItem.class);
}

public void ackWork(String environmentId, String workId, String workerId) {
Expand All @@ -109,17 +110,16 @@ public void ackWork(String environmentId, String workId, String workerId) {
execute(lifecycleApi.ackEnvironmentWork(environmentId, workId, workerHeader(workerId)));
}

public HeartbeatResponse heartbeatWork(
public HeartbeatWorkResponse heartbeatWork(
String environmentId, String workId, String expectedLastHeartbeat, int desiredTTLSeconds) {
require(environmentId, "environmentId");
require(workId, "workId");
String expected = expectedLastHeartbeat == null || expectedLastHeartbeat.isEmpty()
? SelfHostedConstants.EXPECTED_LAST_HEARTBEAT_NO_HEARTBEAT
: expectedLastHeartbeat;
Integer ttl = desiredTTLSeconds > 0 ? desiredTTLSeconds : null;
HeartbeatWorkResponse response = execute(heartbeatApi.heartbeatEnvironmentWork(
return execute(heartbeatApi.heartbeatEnvironmentWork(
environmentId, workId, expected, ttl, Collections.<String, String>emptyMap()));
return HeartbeatResponse.fromMap(toMap(response));
}

public void stopWork(String environmentId, String workId, boolean force) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,6 @@ public final class SelfHostedConstants {
public static final String EVENT_LIST_ORDER_ASC = "asc";
public static final String SESSION_STOP_REASON_END_TURN = "end_turn";

public static final String WORK_STATE_STOPPING = "stopping";
public static final String WORK_STATE_STOPPED = "stopped";

public static final long DEFAULT_MAX_IDLE_MILLIS = 60000L;
public static final long DEFAULT_TOOL_TIMEOUT_MILLIS = 120000L;
public static final long DEFAULT_HEARTBEAT_MILLIS = 30000L;
Expand Down
51 changes: 0 additions & 51 deletions src/main/java/com/volcengine/ark/runtime/selfhosted/WorkData.java

This file was deleted.

Loading
Loading