Skip to content

Commit 5bd07cb

Browse files
committed
feat(agents): expose turn-scoped artifact content responses
1 parent 5541f50 commit 5bd07cb

3 files changed

Lines changed: 113 additions & 26 deletions

File tree

‎docs/helpers/agent-files.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,3 +35,14 @@ var docs = AgentEnvironmentFiles.prepareDirectory(client, Paths.get("docs"),
3535
var staged = AgentEnvironmentFiles.upload(client, environmentId,
3636
Paths.get("update.csv"), "/workspace/update.csv");
3737
```
38+
39+
Artifact content can be read in memory, or streamed to an application-owned, safe destination path. Both modes select the exact result turn and path; async service overloads return `CompletableFuture`.
40+
41+
```java
42+
var artifacts = AgentArtifactDownloads.forResult(
43+
client.beta().agents().sessions().artifacts(), result);
44+
try (var response = artifacts.content("/workspace/outputs/report.md")) {
45+
byte[] report = response.body().readAllBytes();
46+
}
47+
artifacts.download("/workspace/outputs/report.md", Paths.get("downloaded-report.md"));
48+
```

‎openai-java-core/src/main/kotlin/com/openai/helpers/beta/agents/AgentArtifactDownloads.kt‎

Lines changed: 45 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -50,13 +50,30 @@ private constructor(
5050
}
5151
}
5252

53-
/** Streams the exact turn/path artifact to the caller-selected file, replacing that file. */
53+
/** Returns the exact turn/path response. The caller owns and must close it. */
54+
@JvmOverloads
55+
fun content(path: String, options: RequestOptions = RequestOptions.none()): HttpResponse =
56+
openContent(findArtifact(path, options), options)
57+
58+
/** Streams to an application-owned destination, replacing that file. */
5459
@JvmOverloads
5560
fun download(
5661
path: String,
5762
destination: Path,
5863
options: RequestOptions = RequestOptions.none(),
5964
): SessionArtifact {
65+
val artifact = findArtifact(path, options)
66+
openContent(artifact, options).use { Files.copy(it.body(), destination, REPLACE_EXISTING) }
67+
return artifact
68+
}
69+
70+
private fun openContent(artifact: SessionArtifact, options: RequestOptions) =
71+
service.content(
72+
ArtifactContentParams.builder().sessionId(sessionId).artifactId(artifact.id()).build(),
73+
options,
74+
)
75+
76+
private fun findArtifact(path: String, options: RequestOptions): SessionArtifact {
6077
var selected: SessionArtifact? = null
6178
var params = ArtifactListParams.builder().sessionId(sessionId).build()
6279
while (true) {
@@ -65,17 +82,7 @@ private constructor(
6582
if (!page.hasNextPage()) break
6683
params = page.nextPageParams()
6784
}
68-
val artifact = checkNotNull(selected) { "No artifact matches this turn and path" }
69-
service
70-
.content(
71-
ArtifactContentParams.builder()
72-
.sessionId(sessionId)
73-
.artifactId(artifact.id())
74-
.build(),
75-
options,
76-
)
77-
.use { Files.copy(it.body(), destination, REPLACE_EXISTING) }
78-
return artifact
85+
return checkNotNull(selected) { "No artifact matches this turn and path" }
7986
}
8087

8188
class Async
@@ -84,14 +91,34 @@ private constructor(
8491
private val sessionId: String,
8592
private val turnId: String,
8693
) {
94+
/** Returns the exact turn/path response. The caller owns and must close it. */
95+
@JvmOverloads
96+
fun content(
97+
path: String,
98+
options: RequestOptions = RequestOptions.none(),
99+
): CompletableFuture<HttpResponse> = withContent(path, options) { _, response -> response }
100+
101+
/** Streams to an application-owned destination, replacing that file. */
87102
@JvmOverloads
88103
fun download(
89104
path: String,
90105
destination: Path,
91106
options: RequestOptions = RequestOptions.none(),
92-
): CompletableFuture<SessionArtifact> {
93-
val result = CompletableFuture<SessionArtifact>()
107+
): CompletableFuture<SessionArtifact> =
108+
withContent(path, options) { artifact, response ->
109+
response.use { Files.copy(it.body(), destination, REPLACE_EXISTING) }
110+
artifact
111+
}
112+
113+
private fun <T> withContent(
114+
path: String,
115+
options: RequestOptions,
116+
consume: (SessionArtifact, HttpResponse) -> T,
117+
): CompletableFuture<T> {
118+
val result = CompletableFuture<T>()
94119
val response = AtomicReference<HttpResponse?>()
120+
// Native service futures are dependent stages: canceling them does not abort their
121+
// transport and can suppress delivery of a response that still needs closing.
95122
result.whenComplete { _, _ -> if (result.isCancelled) response.get()?.close() }
96123
var selected: SessionArtifact? = null
97124
fun downloadSelected() {
@@ -110,12 +137,11 @@ private constructor(
110137
else {
111138
response.set(content)
112139
try {
113-
content.use {
114-
if (!result.isDone)
115-
Files.copy(it.body(), destination, REPLACE_EXISTING)
116-
}
117-
result.complete(artifact)
140+
if (result.isDone) content.close()
141+
else if (!result.complete(consume(artifact, content)))
142+
content.close()
118143
} catch (error: Throwable) {
144+
runCatching { content.close() }
119145
result.completeExceptionally(error)
120146
} finally {
121147
response.set(null)

‎openai-java-core/src/test/kotlin/com/openai/helpers/beta/agents/AgentFileHelpersTest.kt‎

Lines changed: 57 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -440,6 +440,51 @@ internal class AgentFileHelpersTest {
440440
}
441441
}
442442

443+
@ParameterizedTest
444+
@ValueSource(booleans = [false, true])
445+
fun `content preserves exact scope options and caller ownership`(async: Boolean) {
446+
val t =
447+
Transport().apply {
448+
pages =
449+
mutableListOf(
450+
page(artifact("old", "old"), more = true),
451+
page(artifact("selected")),
452+
)
453+
}
454+
val client = t.client()
455+
try {
456+
val response =
457+
if (async)
458+
AgentArtifactDownloads.forResult(
459+
client.async().beta().agents().sessions().artifacts(),
460+
result(),
461+
)
462+
.content("/workspace/outputs/report.txt", options)
463+
.join()
464+
else
465+
AgentArtifactDownloads.forResult(
466+
client.beta().agents().sessions().artifacts(),
467+
result(),
468+
)
469+
.content("/workspace/outputs/report.txt", options)
470+
assertThat(t.contentClosed).isFalse()
471+
assertThat(t.contentReads).isZero()
472+
response.use {
473+
val bytes = ByteArrayOutputStream()
474+
it.body().copyTo(bytes)
475+
assertThat(bytes.size()).isEqualTo(t.contentSize)
476+
}
477+
assertThat(t.contentClosed).isTrue()
478+
assertThat(t.requests).hasSize(3)
479+
assertThat(t.requests.last().pathSegments).contains("s", "selected")
480+
assertThat(t.requestOptions).allSatisfy {
481+
assertThat(it.timeout).isEqualTo(options.timeout)
482+
}
483+
} finally {
484+
client.close()
485+
}
486+
}
487+
443488
@ParameterizedTest
444489
@ValueSource(booleans = [false, true])
445490
fun `artifact lookup reports missing and ambiguous matches without touching destination`(
@@ -597,8 +642,11 @@ internal class AgentFileHelpersTest {
597642
}
598643
}
599644

600-
@Test
601-
fun `cancelling an async download closes a late content response without writing`() {
645+
@ParameterizedTest
646+
@ValueSource(booleans = [false, true])
647+
fun `cancelling native async content closes a late response without writing`(
648+
inMemory: Boolean
649+
) {
602650
val t = Transport()
603651
val client = t.client()
604652
val pending = CompletableFuture<HttpResponse>()
@@ -613,12 +661,14 @@ internal class AgentFileHelpersTest {
613661
}
614662
try {
615663
val destination = directory.resolve("cancelled.txt")
616-
val download =
664+
val scoped =
617665
AgentArtifactDownloads.forResult(
618-
client.async().beta().agents().sessions().artifacts(),
619-
result(),
620-
)
621-
.download("/workspace/outputs/report.txt", destination)
666+
client.async().beta().agents().sessions().artifacts(),
667+
result(),
668+
)
669+
val download =
670+
if (inMemory) scoped.content("/workspace/outputs/report.txt")
671+
else scoped.download("/workspace/outputs/report.txt", destination)
622672
assertThat(requested.await(5, TimeUnit.SECONDS)).isTrue()
623673
assertThat(download.cancel(true)).isTrue()
624674
pending.complete(t.execute(requireNotNull(request), options))

0 commit comments

Comments
 (0)