From a8de39a973b340049f69f1cb02f24e0075a343c2 Mon Sep 17 00:00:00 2001 From: "ark-hand[bot]" <315378070+ark-hand[bot]@users.noreply.github.com> Date: Wed, 26 Aug 2026 13:35:27 +0000 Subject: [PATCH] refactor(selfHosted): use generated work models MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 简述 让 self-hosted worker 直接复用 ark-apis 生成的 WorkItem、WorkData 和 HeartbeatWorkResponse,移除重复协议结构。 ## 改动 - 删除 selfhosted 下重复的 WorkItem、WorkData、HeartbeatResponse - SelfHostedClient、WorkPoller 改用 generated environment model - EnvironmentWorker 使用私有 ClaimedWork 运行上下文,不伪造 generated WorkItem - SessionSnapshot/Event/SkillRef 继续作为 runner 的归一化运行时视图 ## 验证 - mvn -q clean test(40 passed) - mvn -q checkstyle:check - STG Java SkillHub + private skill E2E,通过 Python/Java 命令交互 See merge request: !87 Sync-Source-Commit: 3d729e63c5df1d7d26146e7657001fdc476afbd0 Ark-APIs-Commit: 1ae6acad4ff5d332748b13c47c120ecb01f1f185 Hand-Written-Reason: Self-hosted worker model consolidation using generated work contracts through 1ae6acad. Release-Version: 0.4.0 --- .../runtime/selfhosted/EnvironmentWorker.java | 114 ++++++++++-------- .../runtime/selfhosted/HeartbeatResponse.java | 60 --------- .../runtime/selfhosted/SelfHostedClient.java | 8 +- .../selfhosted/SelfHostedConstants.java | 3 - .../ark/runtime/selfhosted/WorkData.java | 51 -------- .../ark/runtime/selfhosted/WorkItem.java | 92 -------------- .../ark/runtime/selfhosted/WorkItems.java | 31 +++++ .../ark/runtime/selfhosted/WorkPoller.java | 5 +- .../selfhosted/EnvironmentWorkerTest.java | 60 ++++++++- .../selfhosted/SelfHostedClientTest.java | 10 +- 10 files changed, 169 insertions(+), 265 deletions(-) delete mode 100644 src/main/java/com/volcengine/ark/runtime/selfhosted/HeartbeatResponse.java delete mode 100644 src/main/java/com/volcengine/ark/runtime/selfhosted/WorkData.java delete mode 100644 src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItem.java create mode 100644 src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItems.java diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorker.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorker.java index 2fbb515..fa7b997 100644 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorker.java +++ b/src/main/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorker.java @@ -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; @@ -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); @@ -74,35 +77,28 @@ 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 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; } @@ -110,7 +106,7 @@ private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOExc 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()) { @@ -118,8 +114,8 @@ private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOExc } 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) @@ -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); @@ -158,19 +154,19 @@ private void handleItem(WorkItem item, boolean useWorkdirAsSession) throws IOExc } } - private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference cause) { + private void heartbeatLoop(ClaimedWork work, AtomicBoolean stop, AtomicReference 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) { @@ -180,21 +176,13 @@ private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference 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"); @@ -222,8 +219,8 @@ private void heartbeatLoop(WorkItem item, AtomicBoolean stop, AtomicReference raw) { - HeartbeatResponse response = new HeartbeatResponse(); - if (raw == null) { - return response; - } - response.lastHeartbeat = stringValue(raw.get("last_heartbeat")); - if (raw.get("lease_extended") instanceof Boolean) { - response.leaseExtended = (Boolean) raw.get("lease_extended"); - } - response.state = stringValue(raw.get("state")); - response.ttlSeconds = intValue(raw.get("ttl_seconds")); - response.type = stringValue(raw.get("type")); - return response; - } - - private static String stringValue(Object value) { - return value == null ? "" : String.valueOf(value); - } - - private static int intValue(Object value) { - if (value instanceof Number) { - return ((Number) value).intValue(); - } - return 0; - } - - public String getLastHeartbeat() { - return lastHeartbeat; - } - - public Boolean getLeaseExtended() { - return leaseExtended; - } - - public String getState() { - return state; - } - - public int getTtlSeconds() { - return ttlSeconds; - } - - public String getType() { - return type; - } -} diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClient.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClient.java index 8cebce8..617a3f3 100644 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClient.java +++ b/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClient.java @@ -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; @@ -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) { @@ -109,7 +110,7 @@ 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"); @@ -117,9 +118,8 @@ public HeartbeatResponse heartbeatWork( ? 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.emptyMap())); - return HeartbeatResponse.fromMap(toMap(response)); } public void stopWork(String environmentId, String workId, boolean force) { diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedConstants.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedConstants.java index cffa3e9..114f791 100644 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedConstants.java +++ b/src/main/java/com/volcengine/ark/runtime/selfhosted/SelfHostedConstants.java @@ -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; diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkData.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkData.java deleted file mode 100644 index ca07a2a..0000000 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkData.java +++ /dev/null @@ -1,51 +0,0 @@ -// Copyright (c) 2026 ByteDance Ltd. and/or its affiliates. -// SPDX-License-Identifier: Apache-2.0 - -package com.volcengine.ark.runtime.selfhosted; - -import java.util.Map; - -public class WorkData { - private String type = ""; - private String id = ""; - private String sessionId = ""; - - public static WorkData fromMap(Map raw) { - WorkData data = new WorkData(); - if (raw == null) { - return data; - } - data.type = stringValue(raw.get("type")); - data.id = stringValue(raw.get("id")); - data.sessionId = stringValue(raw.get("session_id")); - return data; - } - - private static String stringValue(Object value) { - return value == null ? "" : String.valueOf(value); - } - - public String getType() { - return type; - } - - public void setType(String type) { - this.type = type; - } - - public String getId() { - return id; - } - - public void setId(String id) { - this.id = id; - } - - public String getSessionId() { - return sessionId; - } - - public void setSessionId(String sessionId) { - this.sessionId = sessionId; - } -} diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItem.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItem.java deleted file mode 100644 index 2fea60f..0000000 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItem.java +++ /dev/null @@ -1,92 +0,0 @@ -// Copyright (c) 2026 ByteDance Ltd. and/or its affiliates. -// SPDX-License-Identifier: Apache-2.0 - -package com.volcengine.ark.runtime.selfhosted; - -import java.util.Map; - -public class WorkItem { - private String id = ""; - private String environmentId = ""; - private WorkData data = new WorkData(); - private String latestHeartbeatAt = ""; - private String sessionId = ""; - private String state = ""; - private String lastHeartbeat = ""; - - @SuppressWarnings("unchecked") - public static WorkItem fromMap(Map raw) { - WorkItem item = new WorkItem(); - if (raw == null) { - return item; - } - item.id = stringValue(raw.get("id")); - item.environmentId = stringValue(raw.get("environment_id")); - item.latestHeartbeatAt = stringValue(raw.get("latest_heartbeat_at")); - item.sessionId = stringValue(raw.get("session_id")); - item.state = stringValue(raw.get("state")); - item.lastHeartbeat = stringValue(raw.get("last_heartbeat")); - if (raw.get("data") instanceof Map) { - item.data = WorkData.fromMap((Map) raw.get("data")); - } - return item; - } - - public String sessionIdValue() { - if (sessionId != null && !sessionId.isEmpty()) { - return sessionId; - } - if (data.getSessionId() != null && !data.getSessionId().isEmpty()) { - return data.getSessionId(); - } - if (data.getId() != null && !data.getId().isEmpty() - && (data.getType() == null || data.getType().isEmpty() || "session".equals(data.getType()))) { - return data.getId(); - } - return ""; - } - - public String latestHeartbeatValue() { - return latestHeartbeatAt != null && !latestHeartbeatAt.isEmpty() ? latestHeartbeatAt : lastHeartbeat; - } - - private static String stringValue(Object value) { - return value == null ? "" : String.valueOf(value); - } - - public String getId() { - return id; - } - - public void setId(String id) { - this.id = id; - } - - public String getEnvironmentId() { - return environmentId; - } - - public void setEnvironmentId(String environmentId) { - this.environmentId = environmentId; - } - - public WorkData getData() { - return data; - } - - public void setData(WorkData data) { - this.data = data; - } - - public String getLatestHeartbeatAt() { - return latestHeartbeatAt; - } - - public void setLatestHeartbeatAt(String latestHeartbeatAt) { - this.latestHeartbeatAt = latestHeartbeatAt; - } - - public String getState() { - return state; - } -} diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItems.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItems.java new file mode 100644 index 0000000..db303d8 --- /dev/null +++ b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkItems.java @@ -0,0 +1,31 @@ +// Copyright (c) 2026 ByteDance Ltd. and/or its affiliates. +// SPDX-License-Identifier: Apache-2.0 + +package com.volcengine.ark.runtime.selfhosted; + +import com.volcengine.ark.runtime.models.environment.WorkData; +import com.volcengine.ark.runtime.models.environment.WorkItem; + +final class WorkItems { + private WorkItems() { + } + + static String sessionId(WorkItem item) { + if (item == null) { + return ""; + } + WorkData data = item.getData(); + if (data == null || data.getId() == null || data.getId().isEmpty()) { + return ""; + } + String type = data.getType(); + return type == null || type.isEmpty() || "session".equals(type) ? data.getId() : ""; + } + + static String latestHeartbeat(WorkItem item) { + if (item == null || item.getLatestHeartbeatAt() == null) { + return ""; + } + return item.getLatestHeartbeatAt(); + } +} diff --git a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkPoller.java b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkPoller.java index ee5a5f9..857a830 100644 --- a/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkPoller.java +++ b/src/main/java/com/volcengine/ark/runtime/selfhosted/WorkPoller.java @@ -3,6 +3,7 @@ package com.volcengine.ark.runtime.selfhosted; +import com.volcengine.ark.runtime.models.environment.WorkItem; import java.util.Random; import java.util.logging.Level; import java.util.logging.Logger; @@ -66,7 +67,7 @@ public WorkItem next() { if (item.getEnvironmentId() == null || item.getEnvironmentId().isEmpty()) { item.setEnvironmentId(options.environmentId); } - if (item.sessionIdValue().isEmpty()) { + if (WorkItems.sessionId(item).isEmpty()) { options.logger.warning( "discard invalid work work_id=" + item.getId() + " reason=missing session id"); discardInvalidWork(item); @@ -91,7 +92,7 @@ public WorkItem next() { pendingStop = () -> stopItem(item, false); } discards = 0; - options.logger.info("claimed work work_id=" + item.getId() + " session_id=" + item.sessionIdValue()); + options.logger.info("claimed work work_id=" + item.getId() + " session_id=" + WorkItems.sessionId(item)); return item; } return null; diff --git a/src/test/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorkerTest.java b/src/test/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorkerTest.java index 5d4fd2c..e4dc1b7 100644 --- a/src/test/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorkerTest.java +++ b/src/test/java/com/volcengine/ark/runtime/selfhosted/EnvironmentWorkerTest.java @@ -6,6 +6,8 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +import com.volcengine.ark.runtime.models.environment.HeartbeatWorkResponse; +import com.volcengine.ark.runtime.models.environment.WorkState; import java.io.IOException; import java.nio.file.Files; import java.util.concurrent.CountDownLatch; @@ -49,7 +51,7 @@ public void heartbeatStopCancelsRunnerBeforeEventPolling() throws Exception { String path = request.url().encodedPath(); if (path.endsWith("/heartbeat")) { heartbeat.countDown(); - return response(request, "{\"state\":\"stopping\",\"lease_extended\":true,\"ttl_seconds\":30}"); + return response(request, "{\"state\":\"stopping\",\"lease_extended\":true}"); } if (path.endsWith("/sessions/session-1")) { try { @@ -105,6 +107,22 @@ public void leaseLostDoesNotStopWorkOwnedByAnotherWorker() throws Exception { assertEquals(0, client.stops.get()); } + @Test + public void leaseNotExtendedWithoutTtlDoesNotStopWorkOwnedByAnotherWorker() throws Exception { + LeaseNotExtendedClient client = new LeaseNotExtendedClient(); + EnvironmentWorker worker = new EnvironmentWorker( + client, + new EnvironmentWorker.Options() + .workdir(Files.createTempDirectory("ark-java-worker-").toString())); + + worker.handleItem(new EnvironmentWorker.HandleItemOptions() + .environmentId("env-1") + .workId("work-1") + .sessionId("session-1")); + + assertEquals(0, client.stops.get()); + } + private static Response response(Request request, String body) throws IOException { return new Response.Builder() .request(request) @@ -124,7 +142,7 @@ private static class NullHeartbeatClient extends SelfHostedClient { } @Override - public HeartbeatResponse heartbeatWork( + public HeartbeatWorkResponse heartbeatWork( String environmentId, String workId, String expectedLastHeartbeat, int desiredTTLSeconds) { heartbeats.incrementAndGet(); firstHeartbeat.countDown(); @@ -157,7 +175,7 @@ private static class LeaseLostClient extends SelfHostedClient { } @Override - public HeartbeatResponse heartbeatWork( + public HeartbeatWorkResponse heartbeatWork( String environmentId, String workId, String expectedLastHeartbeat, int desiredTTLSeconds) { heartbeat.countDown(); throw new WorkerAPIException(412, "lease lost", ""); @@ -181,4 +199,40 @@ public void stopWork(String environmentId, String workId, boolean force) { stops.incrementAndGet(); } } + + private static class LeaseNotExtendedClient extends SelfHostedClient { + private final CountDownLatch heartbeat = new CountDownLatch(1); + private final AtomicInteger stops = new AtomicInteger(); + + LeaseNotExtendedClient() { + super("test-key"); + } + + @Override + public HeartbeatWorkResponse heartbeatWork( + String environmentId, String workId, String expectedLastHeartbeat, int desiredTTLSeconds) { + heartbeat.countDown(); + return new HeartbeatWorkResponse() + .state(WorkState.ACTIVE) + .leaseExtended(Boolean.FALSE); + } + + @Override + public SessionSnapshot getSession(String sessionId) { + try { + assertTrue(heartbeat.await(2, TimeUnit.SECONDS)); + } catch (InterruptedException error) { + Thread.currentThread().interrupt(); + throw new RuntimeException(error); + } + SessionSnapshot session = new SessionSnapshot(); + session.setId(sessionId); + return session; + } + + @Override + public void stopWork(String environmentId, String workId, boolean force) { + stops.incrementAndGet(); + } + } } diff --git a/src/test/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClientTest.java b/src/test/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClientTest.java index cb38813..a081f6e 100644 --- a/src/test/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClientTest.java +++ b/src/test/java/com/volcengine/ark/runtime/selfhosted/SelfHostedClientTest.java @@ -10,6 +10,8 @@ import com.sun.net.httpserver.HttpExchange; import com.sun.net.httpserver.HttpServer; import com.volcengine.ark.runtime.interceptor.RetryInterceptor; +import com.volcengine.ark.runtime.models.environment.WorkItem; +import com.volcengine.ark.runtime.models.environment.WorkState; import java.io.IOException; import java.net.InetSocketAddress; import java.nio.charset.StandardCharsets; @@ -61,7 +63,7 @@ public void preservesNestedSessionWorkData() { assertEquals("work-1", item.getId()); assertEquals("env-1", item.getEnvironmentId()); assertEquals("session-1", item.getData().getId()); - assertEquals("session-1", item.sessionIdValue()); + assertEquals("session-1", WorkItems.sessionId(item)); } @Test @@ -255,7 +257,11 @@ public void atomicEnvironmentWorkAPIMatchesOpenAPIContract() throws Exception { .httpClient(httpClient) .build(); - assertEquals("work-1", client.pollWork("env-1", "worker-1", 999, 5000).getId()); + WorkItem item = client.pollWork("env-1", "worker-1", 999, 5000); + assertEquals("work-1", item.getId()); + assertEquals("2026-08-24T10:00:00Z", item.getCreatedAt()); + assertEquals(WorkState.ACTIVE, item.getState()); + assertEquals(WorkItem.TypeEnum.WORK, item.getType()); client.ackWork("env-1", "work-1", "worker-1"); assertEquals( "2026-08-24T10:00:01Z",