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
14 changes: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,17 @@ libraryDependencies += "com.tjclp" %%% "scalagent" % "0.11.0"
- Bun (preferred) or Node.js 18+
- `bun install` to fetch the TypeScript SDK and ZIO dependencies

### 0.13.1 Notes

- `GetTask` accepts `includeArtifacts` (JSON-RPC/REST `includeArtifacts` /
`include_artifacts`, gRPC `optional bool include_artifacts = 4`). Absent or
true keeps the spec-shaped response with artifacts; false trims them for
lean status polls. The default is the opposite of ListTasks (which excludes
unless asked), and suppressed artifacts are omitted from the wire — a
trimmed result is indistinguishable from a task with no artifacts and must
not be written back to a store. Local protocol extension, recorded in
`proto/A2A_PROTO_SOURCE.txt`.

### 0.11.0 Notes

- Bump the Claude Agent SDK baseline to `^0.3.201` (from `^0.3.156`) and the Codex SDK to `^0.142.5` (from `^0.134.0`), with facade parity for the new SDK surface: observer agents, `promptId` on hook inputs, sandbox credential protection, per-MCP-server permission-mode overrides, `reinitialize`, thinking display control, and typed `informational` / `model_refusal_*` / `worker_shutting_down` system events. See `docs/COMPATIBILITY.md` for the full matrix.
Expand Down Expand Up @@ -348,6 +359,9 @@ The v1 agent card advertises `JSONRPC` first and `HTTP+JSON` second, both with
`protocolVersion = "1.0"`. `A2AClient.discover` fetches the well-known card and
selects a JSON-RPC endpoint.

`GetTask` accepts `historyLength` and `includeArtifacts` (default include —
the opposite of ListTasks; pass `false` for lean status polls).

`SendMessage` defaults to `returnImmediately = false`, so it waits until the
task reaches a terminal or interrupted state. Use `client.submit(...)` or set
`MessageSendConfiguration(returnImmediately = true)` to get the initial
Expand Down
2 changes: 1 addition & 1 deletion build.mill
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ private def sharedPomSettings = PomSettings(
)

private def sharedPublishVersion = Task.Input {
Task.env.get("PUBLISH_VERSION").getOrElse("0.12.0")
Task.env.get("PUBLISH_VERSION").getOrElse("0.13.1")
}

// Empty Javadoc jar (Scala.js facades trip Scaladoc; we keep parity on JVM).
Expand Down
5 changes: 5 additions & 0 deletions proto/A2A_PROTO_SOURCE.txt
Original file line number Diff line number Diff line change
@@ -1,2 +1,7 @@
a2a.proto vendored from github.com/a2aproject/A2A specification/a2a.proto
upstream commit: cd87b9341bc0e4982d46550aaab2319b903271e4

Local deltas (re-apply after re-vendoring; A2AProtoParitySpec fails until they are):
- `tenant` field on every request message (multi-tenant routing)
- GetTaskRequest: `optional bool include_artifacts = 4` (lean status polls;
tag collides if upstream ever adds a field 4 — resolve at re-vendor time)
5 changes: 5 additions & 0 deletions proto/a2a.proto
Original file line number Diff line number Diff line change
Expand Up @@ -669,6 +669,11 @@ message GetTaskRequest {
// a request to not include any messages. The server MUST NOT return more
// messages than the provided value, but MAY apply a lower limit.
optional int32 history_length = 3;
// Whether to include artifacts in the returned task. Unset defaults to
// true, keeping the spec-shaped GetTask response; false trims the
// artifacts array for lean status polls. Local extension (the upstream
// spec defines this switch only on ListTasksRequest), like `tenant`.
optional bool include_artifacts = 4;
}

// Parameters for listing tasks with optional filtering criteria.
Expand Down
7 changes: 4 additions & 3 deletions src/shared/com/tjclp/scalagent/a2a/A2APathRouting.scala
Original file line number Diff line number Diff line change
Expand Up @@ -171,9 +171,10 @@ private[a2a] object A2APathRouting:
tenant: Option[String],
): Either[A2AError, A2ARequest.TasksGet] =
for
id <- taskId(rawTaskId)
historyLength <- query.int("historyLength", "history_length")
yield A2ARequest.TasksGet(id, historyLength, tenant)
id <- taskId(rawTaskId)
historyLength <- query.int("historyLength", "history_length")
includeArtifacts <- query.bool("includeArtifacts", "include_artifacts")
yield A2ARequest.TasksGet(id, historyLength, includeArtifacts, tenant)

def tasksCancel(
rawTaskAction: String,
Expand Down
22 changes: 17 additions & 5 deletions src/shared/com/tjclp/scalagent/a2a/A2ARequest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -141,10 +141,19 @@ object A2ARequest:
type SendMessageRequest = MessageSend
type SendStreamingMessageRequest = MessageSend

/** Parameters for GetTask. */
/**
* Parameters for GetTask.
*
* `includeArtifacts`: absent or true returns artifacts (the spec-shaped
* response); false trims them for lean status polls — a local extension
* with the opposite default of ListTasks. Suppressed artifacts are
* omitted from the wire, indistinguishable from a task that has none, so
* a trimmed GetTask result must never be written back to a store.
*/
final case class TasksGet(
id: TaskId,
historyLength: Option[Int] = None,
includeArtifacts: Option[Boolean] = None,
tenant: Option[String] = None)
object TasksGet:
given JsonEncoder[TasksGet] = JsonEncoder[Json].contramap { request =>
Expand All @@ -153,17 +162,20 @@ object A2ARequest:
request.historyLength.foreach(value =>
obj = obj.add("historyLength", Json.Num(java.math.BigDecimal.valueOf(value.toLong)))
)
request.includeArtifacts.foreach(value => obj = obj.add("includeArtifacts", Json.Bool(value)))
obj
}
given JsonDecoder[TasksGet] = JsonDecoder[Json].mapOrFail { json =>
objectFields(json, "GetTaskRequest").flatMap { fields =>
for
id <- requiredString(fields, "id")
historyLength <- optionalNonNegativeInt(fields, "historyLength", "history_length")
tenant <- optionalString(fields, "tenant")
yield TasksGet(TaskId(id), historyLength, tenant)
id <- requiredString(fields, "id")
historyLength <- optionalNonNegativeInt(fields, "historyLength", "history_length")
includeArtifacts <- optionalBool(fields, "includeArtifacts", "include_artifacts")
tenant <- optionalString(fields, "tenant")
yield TasksGet(TaskId(id), historyLength, includeArtifacts, tenant)
}
}
end TasksGet
type GetTaskRequest = TasksGet

/** Parameters for ListTasks. */
Expand Down
11 changes: 10 additions & 1 deletion src/shared/com/tjclp/scalagent/a2a/A2ARequestHandler.scala
Original file line number Diff line number Diff line change
Expand Up @@ -127,11 +127,20 @@ private[a2a] final class A2ARequestHandler(
}
}

// GetTask includes artifacts unless the request says includeArtifacts=false
// (absent means include — the opposite of ListTasks, whose spec default is
// exclude; the spec defines no artifact switch for GetTask, so this field
// is a local extension and absent must keep the spec-conformant response).
def getTask(params: A2ARequest.TasksGet, context: ServerCallContext): Task[A2ATask] =
authorizeRequest(context) *> validateHistoryLength(params.historyLength) *>
taskStore.load(params.id, context.tenant).flatMap {
case Some(task) =>
reconcileOrphaned(task, context).map(A2ATaskStore.applyHistoryLength(_, params.historyLength))
reconcileOrphaned(task, context).map { reconciled =>
A2ATaskStore.applyArtifactProjection(
A2ATaskStore.applyHistoryLength(reconciled, params.historyLength),
includeArtifacts = params.includeArtifacts.getOrElse(true),
)
}
case None => ZIO.fail(A2AError.taskNotFound(params.id))
}

Expand Down
13 changes: 11 additions & 2 deletions src/shared/com/tjclp/scalagent/a2a/A2AServerTypes.scala
Original file line number Diff line number Diff line change
Expand Up @@ -464,15 +464,24 @@ object A2ATaskStore:
val includeArtifacts = params.includeArtifacts.getOrElse(false)
A2AResponse.ListTasksResult(
tasks = page.map { task =>
val withHistory = applyHistoryLength(task, params.historyLength)
if includeArtifacts then withHistory else withHistory.copy(artifacts = Nil)
applyArtifactProjection(applyHistoryLength(task, params.historyLength), includeArtifacts)
},
nextPageToken = next,
pageSize = pageSize,
totalSize = filtered.length,
includeArtifacts = includeArtifacts,
)

/**
* Drop `task.artifacts` unless `includeArtifacts` asks for them. The
* caller states its default explicitly: GetTask includes on absent (the
* spec-shaped response), ListTasks excludes on absent (the spec's
* payload-size default). Public for the same reason as
* [[applyHistoryLength]].
*/
def applyArtifactProjection(task: A2ATask, includeArtifacts: Boolean): A2ATask =
if includeArtifacts then task else task.copy(artifacts = Nil)

/**
* Truncate `task.history` to the requested length (matches the A2A
* `historyLength` semantic). Public so durable [[A2ATaskStore]] impls
Expand Down
5 changes: 5 additions & 0 deletions test/jvm/com/tjclp/scalagent/a2a/A2AGrpcProtoCodecSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@ class A2AGrpcProtoCodecSpec extends FunSuite:
.setTenant("tenant-a")
.setId("task-1")
.setHistoryLength(2)
// false is the value presence-tracking protects: without `optional` in
// the vendored proto, JsonPrinter drops the proto3 default and the
// lean poll silently becomes a full one.
.setIncludeArtifacts(false)
.build()

val decoded = A2AGrpcProtoCodec.decodeRequest(A2AOperation.TasksGet, proto.toByteArray)
Expand All @@ -25,6 +29,7 @@ class A2AGrpcProtoCodecSpec extends FunSuite:
assertEquals(request.id, TaskId("task-1"))
assertEquals(request.tenant, Some("tenant-a"))
assertEquals(request.historyLength, Some(2))
assertEquals(request.includeArtifacts, Some(false))
case other =>
fail(s"expected decoded TasksGet request, got $other")

Expand Down
6 changes: 4 additions & 2 deletions test/jvm/com/tjclp/scalagent/a2a/A2AProtoParitySpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1019,7 +1019,7 @@ class A2AProtoParitySpec extends FunSuite:

private def expectedRestQueryFields: Map[String, Set[String]] =
Map(
A2AMethod.TasksGet -> Set("historyLength"),
A2AMethod.TasksGet -> Set("historyLength", "includeArtifacts"),
A2AMethod.TasksList -> Set(
"contextId",
"status",
Expand Down Expand Up @@ -1251,7 +1251,9 @@ class A2AProtoParitySpec extends FunSuite:
"SendMessageRequest" -> jsonFields(
A2ARequest.MessageSend(message, configuration = Some(sendConfig), metadata = Some(Json.Obj("m" -> Json.Str("request"))), tenant = Some("tenant-a"))
),
"GetTaskRequest" -> jsonFields(A2ARequest.TasksGet(taskId, historyLength = Some(2), tenant = Some("tenant-a"))),
"GetTaskRequest" -> jsonFields(
A2ARequest.TasksGet(taskId, historyLength = Some(2), includeArtifacts = Some(false), tenant = Some("tenant-a"))
),
"ListTasksRequest" -> jsonFields(
A2ARequest.TasksList(
contextId = Some(contextId),
Expand Down
54 changes: 54 additions & 0 deletions test/jvm/com/tjclp/scalagent/a2a/A2AServerLiveSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -600,6 +600,60 @@ class A2AServerLiveSpec extends FunSuite:
assert(wrong.error.exists(_.message.contains("selected AgentInterface tenant")))
}

test("JVM REST and JSON-RPC GetTask honor includeArtifacts=false"):
val program =
for
taskStore <- ZIO.succeed(A2ATaskStore.inMemory)
taskId = TaskId("artifact-wire-task")
contextId = ContextId("artifact-wire-context")
task = A2ATask(
id = taskId,
contextId = contextId,
status = TaskStatus.completed(A2AMessage.agentText("done", Some(contextId))),
artifacts = List(Artifact("wire-artifact", parts = List(Part.Text("payload")))),
)
_ <- taskStore.save(task, None)
server <- testServer(
A2AServerLive.Config(
name = "GetTaskArtifactTrimJvmTest",
description = "GetTask includeArtifacts JVM test server",
taskStore = Some(taskStore),
)
)
fullResponse <- server.handleHttp(
Request
.get(URL.decode(s"/tasks/${taskId.value}").toOption.get)
.copy(headers = versionedJsonHeaders(A2AContentType.A2AJson))
)
fullBody <- fullResponse.body.asString
leanResponse <- server.handleHttp(
Request
.get(URL.decode(s"/tasks/${taskId.value}?include_artifacts=false").toOption.get)
.copy(headers = versionedJsonHeaders(A2AContentType.A2AJson))
)
leanBody <- leanResponse.body.asString
rpcResponse <- server.handleHttp(
Request
.post(
"/",
Body.fromString(
"""{"jsonrpc":"2.0","id":91,"method":"GetTask","params":{"id":"artifact-wire-task","includeArtifacts":false}}"""
),
)
.copy(headers = versionedJsonHeaders())
)
rpcBody <- rpcResponse.body.asString
yield (fullBody, leanBody, rpcBody)

runTask(program).map { case (fullBody, leanBody, rpcBody) =>
// Suppressed artifacts are omitted entirely, not emitted as [] — the
// lean body must not carry the key at all.
assert(fullBody.contains(""""artifacts""""))
assert(!leanBody.contains(""""artifacts""""))
assert(rpcBody.contains(""""result""""))
assert(!rpcBody.contains(""""artifacts""""))
}

test("JVM JSON-RPC distinguishes malformed JSON from invalid request envelopes"):
val invalidEnvelope =
"""{
Expand Down
14 changes: 10 additions & 4 deletions test/shared/com/tjclp/scalagent/a2a/A2ACodecSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1040,7 +1040,7 @@ class A2ACodecSpec extends FunSuite:

test("request and response model groups round-trip"):
roundTrip(A2ARequest.MessageSend(message, Some(MessageSendConfiguration(taskPushNotificationConfig = Some(pushConfig), historyLength = Some(1)))))
roundTrip(A2ARequest.TasksGet(taskId, historyLength = Some(1), tenant = Some("tenant-codec")))
roundTrip(A2ARequest.TasksGet(taskId, historyLength = Some(1), includeArtifacts = Some(false), tenant = Some("tenant-codec")))
roundTrip(A2ARequest.TasksList(contextId = Some(contextId), status = Some(TaskState.Completed), pageSize = Some(10), pageToken = Some("10")))
roundTrip(A2ARequest.TasksCancel(taskId))
roundTrip(A2ARequest.TasksResubscribe(taskId))
Expand All @@ -1055,6 +1055,7 @@ class A2ACodecSpec extends FunSuite:
val bareMessage = message.copy(metadata = None)
val messageSendJson = A2ARequest.MessageSend(bareMessage).toJson
val getJson = A2ARequest.TasksGet(taskId).toJson
val getFalseJson = A2ARequest.TasksGet(taskId, includeArtifacts = Some(false)).toJson
val listEmptyJson = A2ARequest.TasksList().toJson
val listFalseJson = A2ARequest.TasksList(includeArtifacts = Some(false)).toJson
val extendedJson = A2ARequest.GetAuthenticatedExtendedCard().toJson
Expand All @@ -1066,7 +1067,8 @@ class A2ACodecSpec extends FunSuite:
assertEquals(listEmptyJson, "{}")
assertEquals(extendedJson, "{}")
assert(listFalseJson.contains(""""includeArtifacts":false"""))
assert(!List(messageSendJson, getJson, listEmptyJson, listFalseJson, extendedJson).exists(_.contains("null")))
assert(getFalseJson.contains(""""includeArtifacts":false"""))
assert(!List(messageSendJson, getJson, getFalseJson, listEmptyJson, listFalseJson, extendedJson).exists(_.contains("null")))

test("ProtoJSON null optional fields decode as unset"):
val cardJson =
Expand Down Expand Up @@ -1828,7 +1830,8 @@ class A2ACodecSpec extends FunSuite:
val taskGet =
"""{
| "id": "task-snake",
| "history_length": 3
| "history_length": 3,
| "include_artifacts": false
|}""".stripMargin.fromJson[A2ARequest.TasksGet]
val taskList =
"""{
Expand All @@ -1852,7 +1855,10 @@ class A2ACodecSpec extends FunSuite:
| "page_token": "next"
|}""".stripMargin.fromJson[A2ARequest.PushNotificationConfigList]

assertEquals(taskGet.map(value => (value.id, value.historyLength)), Right((TaskId("task-snake"), Some(3))))
assertEquals(
taskGet.map(value => (value.id, value.historyLength, value.includeArtifacts)),
Right((TaskId("task-snake"), Some(3), Some(false))),
)
assertEquals(
taskList.map(value => (value.contextId, value.pageSize, value.historyLength, value.includeArtifacts)),
Right((Some(ContextId("ctx-snake")), Some(25), Some(1), Some(true))),
Expand Down
4 changes: 3 additions & 1 deletion test/shared/com/tjclp/scalagent/a2a/A2APathRoutingSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,9 @@ class A2APathRoutingSpec extends FunSuite:
)
assertEquals(
A2APathRouting.tasksGet("task-1", query, Some("tenant-a")),
Right(A2ARequest.TasksGet(TaskId("task-1"), Some(3), Some("tenant-a"))),
Right(
A2ARequest.TasksGet(TaskId("task-1"), Some(3), includeArtifacts = Some(false), tenant = Some("tenant-a"))
),
)
assertEquals(
A2APathRouting.pushConfigList("task-1", query, Some("tenant-a")),
Expand Down
42 changes: 42 additions & 0 deletions test/shared/com/tjclp/scalagent/a2a/A2AServerCoreSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -334,6 +334,48 @@ class A2AServerCoreSpec extends FunSuite:
assertEquals(stored.history.map(_.text), List("run"))
}

test("getTask trims artifacts only when includeArtifacts is false"):
val effect =
for
taskStore <- ZIO.succeed(A2ATaskStore.inMemory)
registry <- A2ARuntimeRegistry.make
taskId = TaskId("artifact-trim-task")
contextId = ContextId("artifact-trim-context")
task = A2ATask(
id = taskId,
contextId = contextId,
status = TaskStatus.completed(A2AMessage.agentText("done", Some(contextId))),
artifacts = List(Artifact("keep-or-trim", parts = List(Part.Text("payload")))),
)
_ <- taskStore.save(task, None)
core = A2AServerCore.make(
TestConfig(taskStore = Some(taskStore)),
runtime,
registry,
NoopPushPoster,
() => card("ArtifactTrim"),
(_, _) => ZIO.unit,
)
defaulted <- core.requestHandler.getTask(A2ARequest.TasksGet(taskId), ServerCallContext())
explicit <- core.requestHandler.getTask(
A2ARequest.TasksGet(taskId, includeArtifacts = Some(true)),
ServerCallContext(),
)
trimmed <- core.requestHandler.getTask(
A2ARequest.TasksGet(taskId, includeArtifacts = Some(false)),
ServerCallContext(),
)
stored <- taskStore.load(taskId, None)
yield (defaulted, explicit, trimmed, stored)

runTask(effect).map { case (defaulted, explicit, trimmed, stored) =>
assertEquals(defaulted.artifacts.map(_.artifactId), List("keep-or-trim"))
assertEquals(explicit.artifacts.map(_.artifactId), List("keep-or-trim"))
assertEquals(trimmed.artifacts, Nil)
// read-projection only: the stored task keeps its artifacts
assertEquals(stored.map(_.artifacts.map(_.artifactId)), Some(List("keep-or-trim")))
}

test("shared request handler preserves inactive interrupted tasks"):
val effect =
for
Expand Down
Loading