Skip to content
Closed
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
ADK changes
PiperOrigin-RevId: 988228337
  • Loading branch information
kvmilos authored and copybara-github committed Sep 29, 2026
commit bd1b9e194d55f0b1b6c2f08a8d8a1f0def4dac6f
4 changes: 4 additions & 0 deletions .gemini/styleguide.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
# Gemini Code Assist Customization

_Placeholder - customization for Gemini Code Assist in this repository is coming
soon._
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@

jobs:
analyze-new-release-for-adk-docs-updates:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
# Dry-run reads only (this repo + the public docs repo) and skips writes, so
# the built-in GITHUB_TOKEN suffices. For --no-dry-run, use a PAT with write
# access to the docs repo.
Expand All @@ -32,16 +32,16 @@

steps:
- name: Checkout repository
uses: actions/checkout@v6

Check failure on line 35 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 35 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

analyze-releases-for-adk-docs-updates.yml:35: unpinned action reference: action is not pinned to a hash (required by blanket policy)

- name: Set up Java
uses: actions/setup-java@v5

Check failure on line 38 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 38 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

analyze-releases-for-adk-docs-updates.yml:38: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
distribution: temurin
java-version: '17'

- name: Cache Maven packages
uses: actions/cache@v5

Check failure on line 44 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 44 in .github/workflows/analyze-releases-for-adk-docs-updates.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

analyze-releases-for-adk-docs-updates.yml:44: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
path: ~/.m2/repository
key: ${{ runner.os }}-maven-${{ hashFiles('**/pom.xml') }}
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/pr-commit-check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -12,16 +12,16 @@

# Defines the jobs that will run as part of the workflow.
jobs:
check-commit-count:

Check warning on line 15 in .github/workflows/pr-commit-check.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

excessive-permissions

pr-commit-check.yml:15: overly broad permissions: default permissions used due to no permissions: block
# The type of runner that the job will run on. 'ubuntu-latest' is a good default.
runs-on: ubuntu-latest
# The type of runner that the job will run on, pinned to a fixed Ubuntu release.
runs-on: ubuntu-24.04

# The steps that will be executed as part of the job.
steps:
# Step 1: Check out the code
# This action checks out your repository under $GITHUB_WORKSPACE, so your workflow can access it.
- name: Checkout Code
uses: actions/checkout@v6

Check failure on line 24 in .github/workflows/pr-commit-check.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 24 in .github/workflows/pr-commit-check.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

pr-commit-check.yml:24: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
# We need to fetch all commits to accurately count them.
# '0' means fetch all history for all branches and tags.
Expand All @@ -47,7 +47,7 @@
if: steps.count_commits.outputs.commit_count > 1
# If the condition is met, the workflow will exit with a failure status.
run: |
echo "This pull request has ${{ steps.count_commits.outputs.commit_count }} commits."

Check failure on line 50 in .github/workflows/pr-commit-check.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/template-injection

code injection via template expansion: may expand into attacker-controllable code
echo "Please squash them into a single commit before merging."
echo "You can use git rebase -i HEAD~N"
echo "...where N is the number of commits you want to squash together. The PR check conveniently tells you this number! For example, if the check says you have 3 commits, you would run: git rebase -i HEAD~3."
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/pr-title-check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ permissions:

jobs:
check-pr-title:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- name: Validate Conventional Commit title
env:
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/pr-triage-adk-java.yml
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
# this repository (see the sample's README for the full threat model).
name: ADK PR Triaging Agent

on:

Check failure on line 23 in .github/workflows/pr-triage-adk-java.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

dangerous-triggers

pr-triage-adk-java.yml:23: use of fundamentally insecure workflow trigger: pull_request_target is almost always used insecurely
pull_request_target:
types: [opened, reopened, edited]
workflow_dispatch:
Expand All @@ -44,7 +44,7 @@

jobs:
agent-triage-pull-request:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
# Only run on the upstream repo, for newly-opened/reopened/edited PRs or a
# manual dispatch.
if: >-
Expand All @@ -61,10 +61,10 @@
steps:
# Default checkout: the base branch (trusted code), NOT the PR head.
- name: Checkout repository
uses: actions/checkout@v6

Check failure on line 64 in .github/workflows/pr-triage-adk-java.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 64 in .github/workflows/pr-triage-adk-java.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

pr-triage-adk-java.yml:64: unpinned action reference: action is not pinned to a hash (required by blanket policy)

- name: Set up Java
uses: actions/setup-java@v5

Check failure on line 67 in .github/workflows/pr-triage-adk-java.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 67 in .github/workflows/pr-triage-adk-java.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

pr-triage-adk-java.yml:67: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
distribution: temurin
java-version: '17'
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/release-please.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@
name: release-please
jobs:
release-please:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: googleapis/release-please-action@v4

Check failure on line 15 in .github/workflows/release-please.yaml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 15 in .github/workflows/release-please.yaml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

release-please.yaml:15: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
token: ${{ secrets.RELEASE_PLEASE_TOKEN }}
2 changes: 1 addition & 1 deletion .github/workflows/spam-detection-adk-java-issues.yml
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@

jobs:
agent-scan-issues:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
# Only run on the upstream repo, for newly-opened issues, the scheduled
# sweep, or a manual dispatch.
if: >-
Expand All @@ -56,10 +56,10 @@

steps:
- name: Checkout repository
uses: actions/checkout@v6

Check failure on line 59 in .github/workflows/spam-detection-adk-java-issues.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

zizmor/unpinned-uses

unpinned action reference: action is not pinned to a hash (required by blanket policy)

Check failure on line 59 in .github/workflows/spam-detection-adk-java-issues.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

spam-detection-adk-java-issues.yml:59: unpinned action reference: action is not pinned to a hash (required by blanket policy)

- name: Set up Java
uses: actions/setup-java@v5

Check failure on line 62 in .github/workflows/spam-detection-adk-java-issues.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

unpinned-uses

spam-detection-adk-java-issues.yml:62: unpinned action reference: action is not pinned to a hash (required by blanket policy)
with:
distribution: temurin
java-version: '17'
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/stale-adk-java-issues.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ concurrency:

jobs:
agent-audit-stale-issues:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
# Only run on the upstream repo, on the daily schedule or a manual dispatch.
if: >-
github.repository == 'google/adk-java' && (
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/triage-adk-java-issues.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ concurrency:

jobs:
agent-triage-issues:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
# Only run on the upstream repo, for newly-opened issues, the scheduled
# batch sweep, or a manual dispatch.
if: >-
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/validation.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@
cancel-in-progress: true

jobs:
build-modules:

Check warning on line 14 in .github/workflows/validation.yml

View workflow job for this annotation

GitHub Actions / zizmor-output

excessive-permissions

validation.yml:14: overly broad permissions: default permissions used due to no permissions: block
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
timeout-minutes: 30
strategy:
matrix:
Expand Down
4 changes: 4 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
# AGENTS.md

_Placeholder - guidance for AI coding agents working in this repository is coming
soon._
1 change: 1 addition & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
See [AGENTS.md](./AGENTS.md) for project context, commands, and contribution guidelines for AI coding agents.
Original file line number Diff line number Diff line change
Expand Up @@ -317,8 +317,6 @@ synchronized void handleError(String message, Throwable e) {
emitter.tryOnError(new A2AClientError(message, e));
}

// TODO: b/483038527 - The synchronized block might block the thread, we should optimize for
// performance in the future.
synchronized void handleEvent(ClientEvent clientEvent, AgentCard unused) {
// Mark the flow as done if it is already cancelled.
if (!done) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,9 @@
import com.google.adk.utils.Constants;
import com.google.api.core.ApiFuture;
import com.google.api.core.ApiFutures;
import com.google.api.gax.rpc.AlreadyExistsException;
import com.google.cloud.firestore.CollectionReference;
import com.google.cloud.firestore.DocumentReference;
import com.google.cloud.firestore.DocumentSnapshot;
import com.google.cloud.firestore.Firestore;
import com.google.cloud.firestore.Query;
Expand All @@ -48,9 +50,10 @@
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.regex.Matcher;
import javax.annotation.Nullable;
import org.jspecify.annotations.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -73,6 +76,9 @@ public class FirestoreSessionService implements BaseSessionService {
private static final String UPDATE_TIME_KEY = Constants.KEY_UPDATE_TIME;
private static final String TIMESTAMP_KEY = Constants.KEY_TIMESTAMP;

/** Random token each create stores, so a retried create can detect its own write. */
static final String CREATE_TOKEN_KEY = "createToken";

/** Constructor for FirestoreSessionService. */
public FirestoreSessionService(Firestore firestore) {
this.firestore = firestore;
Expand All @@ -96,7 +102,10 @@ public Single<Session> createSession(
return createSession(appName, userId, (Map<String, Object>) state, sessionId);
}

/** Creates a new session in Firestore. */
/**
* Creates a new session in Firestore. Session IDs are unique per user across apps, so creating
* one the user already has under any app fails with {@link SessionException}.
*/
@Override
public Single<Session> createSession(
String appName,
Expand Down Expand Up @@ -139,11 +148,22 @@ public Single<Session> createSession(
sessionData.put(USER_ID_KEY, newSession.userId());
sessionData.put(UPDATE_TIME_KEY, newSession.lastUpdateTime().toString());
sessionData.put(STATE_KEY, newSession.state());

// Asynchronously write to Firestore and wait for the result
ApiFuture<WriteResult> future =
getSessionsCollection(userId).document(resolvedSessionId).set(sessionData);
future.get(); // Block until the write is complete
String createToken = UUID.randomUUID().toString();
sessionData.put(CREATE_TOKEN_KEY, createToken);

// Unlike set(), create() fails if the session already exists instead of replacing it.
DocumentReference sessionDoc = getSessionsCollection(userId).document(resolvedSessionId);
try {
sessionDoc.create(sessionData).get();
} catch (ExecutionException e) {
if (!(e.getCause() instanceof AlreadyExistsException)) {
throw e;
}
// A retry after a lost reply fails on this call's own write; its token means success.
if (!createToken.equals(sessionDoc.get().get().getString(CREATE_TOKEN_KEY))) {
throw new SessionException(SessionException.SESSION_ALREADY_EXISTS, e.getCause());
}
}

return newSession;
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,7 +175,8 @@ void run_withUserInput_createsSessionAndExecutesAgent() throws Exception {
when(mockUserDocRef.collection(anyString())).thenReturn(mockSessionsCollection);
when(mockSessionsCollection.document(anyString())).thenReturn(mockSessionDocRef);

when(mockSessionDocRef.set(anyMap())).thenReturn(ApiFutures.immediateFuture(mockWriteResult));
when(mockSessionDocRef.create(anyMap()))
.thenReturn(ApiFutures.immediateFuture(mockWriteResult));
when(mockSessionDocRef.update(anyMap()))
.thenReturn(ApiFutures.immediateFuture(mockWriteResult));
// Mock the event sub-collection chain
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@
import com.google.adk.events.EventActions;
import com.google.adk.utils.Constants;
import com.google.api.core.ApiFutures;
import com.google.api.gax.grpc.GrpcStatusCode;
import com.google.api.gax.rpc.AlreadyExistsException;
import com.google.api.gax.rpc.PermissionDeniedException;
import com.google.api.gax.rpc.UnavailableException;
import com.google.cloud.firestore.CollectionReference;
import com.google.cloud.firestore.DocumentReference;
import com.google.cloud.firestore.DocumentSnapshot;
Expand All @@ -45,17 +49,20 @@
import com.google.common.collect.ImmutableMap;
import com.google.genai.types.Content;
import com.google.genai.types.Part;
import io.grpc.Status;
import io.reactivex.rxjava3.observers.TestObserver;
import java.time.Instant;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;

Expand Down Expand Up @@ -89,6 +96,7 @@ public class FirestoreSessionServiceTest {
@Mock private QuerySnapshot mockQuerySnapshot;
@Mock private WriteResult mockWriteResult;
@Mock private WriteBatch mockWriteBatch;
@Captor private ArgumentCaptor<Map<String, Object>> sessionDataCaptor;

private FirestoreSessionService sessionService;

Expand Down Expand Up @@ -130,7 +138,7 @@ public void setup() {

// Default mock for writes
lenient()
.when(mockSessionDocRef.set(anyMap()))
.when(mockSessionDocRef.create(anyMap()))
.thenReturn(ApiFutures.immediateFuture(mockWriteResult));
lenient()
.when(mockSessionDocRef.update(anyMap()))
Expand Down Expand Up @@ -307,7 +315,7 @@ void createSession_withSessionId_returnsNewSession() {
assertThat(session.id()).isEqualTo(SESSION_ID);
return true;
});
verify(mockSessionDocRef).set(anyMap());
verify(mockSessionDocRef).create(anyMap());
}

/** Tests that createSession creates a new session with a generated session ID. */
Expand All @@ -334,7 +342,7 @@ void createSession_withNullSessionId_generatesNewId() {
assertThat(session.id()).isNotEmpty();
return true;
});
verify(mockSessionDocRef).set(anyMap());
verify(mockSessionDocRef).create(anyMap());
}

/** Tests that createSession creates a new session with an empty session ID. */
Expand All @@ -358,7 +366,7 @@ void createSession_withEmptySessionId_generatesNewId() {
assertThat(session.id()).isNotEqualTo(" ");
return true;
});
verify(mockSessionDocRef).set(anyMap());
verify(mockSessionDocRef).create(anyMap());
}

/** Tests that createSession creates a new session with an empty state when null state is */
Expand Down Expand Up @@ -391,6 +399,114 @@ void createSession_withNullAppName_throwsNullPointerException() {
.assertError(NullPointerException.class);
}

/** Tests that createSession rejects a session ID that is already taken. */
@Test
void createSession_withSessionIdAlreadyTaken_failsWithSessionException() {
// Arrange
when(mockSessionsCollection.document(SESSION_ID)).thenReturn(mockSessionDocRef);
when(mockSessionDocRef.create(anyMap()))
.thenReturn(ApiFutures.immediateFailedFuture(alreadyExists()));
when(mockSessionDocRef.get()).thenReturn(ApiFutures.immediateFuture(mockSessionSnapshot));
// Another caller's session, or one written before sessions carried a token.
when(mockSessionSnapshot.getString(FirestoreSessionService.CREATE_TOKEN_KEY)).thenReturn(null);

// Act
TestObserver<Session> testObserver =
sessionService.createSession(APP_NAME, USER_ID, null, SESSION_ID).test();

// Assert
testObserver.assertError(
e -> {
assertThat(e).isInstanceOf(SessionException.class);
assertThat(e).hasMessageThat().isEqualTo(SessionException.SESSION_ALREADY_EXISTS);
assertThat(e).hasCauseThat().isInstanceOf(AlreadyExistsException.class);
return true;
});
}

/** Tests that createSession succeeds when a client retry fails on the session's own write. */
@Test
void createSession_whenRetryHitsItsOwnWrite_returnsSession() {
// Arrange
when(mockSessionsCollection.document(SESSION_ID)).thenReturn(mockSessionDocRef);
when(mockSessionDocRef.create(sessionDataCaptor.capture()))
.thenReturn(ApiFutures.immediateFailedFuture(alreadyExists()));
when(mockSessionDocRef.get()).thenReturn(ApiFutures.immediateFuture(mockSessionSnapshot));
when(mockSessionSnapshot.getString(FirestoreSessionService.CREATE_TOKEN_KEY))
.thenAnswer(
unused -> sessionDataCaptor.getValue().get(FirestoreSessionService.CREATE_TOKEN_KEY));

// Act
TestObserver<Session> testObserver =
sessionService.createSession(APP_NAME, USER_ID, null, SESSION_ID).test();

// Assert
testObserver.assertComplete();
testObserver.assertValue(session -> session.id().equals(SESSION_ID));
}

/**
* Tests that a failed read-back after a duplicate is propagated, since it cannot tell whose
* session exists.
*/
@Test
void createSession_whenReadBackFails_propagatesFailure() {
// Arrange
when(mockSessionsCollection.document(SESSION_ID)).thenReturn(mockSessionDocRef);
when(mockSessionDocRef.create(anyMap()))
.thenReturn(ApiFutures.immediateFailedFuture(alreadyExists()));
when(mockSessionDocRef.get())
.thenReturn(
ApiFutures.immediateFailedFuture(
new UnavailableException(
"Service unavailable",
/* cause= */ null,
GrpcStatusCode.of(Status.Code.UNAVAILABLE),
/* retryable= */ true)));

// Act
TestObserver<Session> testObserver =
sessionService.createSession(APP_NAME, USER_ID, null, SESSION_ID).test();

// Assert
testObserver.assertError(
e -> {
assertThat(e).isInstanceOf(ExecutionException.class);
assertThat(e).hasCauseThat().isInstanceOf(UnavailableException.class);
return true;
});
}

/**
* Tests that a write failure unrelated to a duplicate session ID is propagated rather than
* reported as one.
*/
@Test
void createSession_whenWriteFailsForAnotherReason_propagatesFailure() {
// Arrange
when(mockSessionsCollection.document(SESSION_ID)).thenReturn(mockSessionDocRef);
when(mockSessionDocRef.create(anyMap()))
.thenReturn(
ApiFutures.immediateFailedFuture(
new PermissionDeniedException(
"Missing or insufficient permissions",
/* cause= */ null,
GrpcStatusCode.of(Status.Code.PERMISSION_DENIED),
/* retryable= */ false)));

// Act
TestObserver<Session> testObserver =
sessionService.createSession(APP_NAME, USER_ID, null, SESSION_ID).test();

// Assert
testObserver.assertError(
e -> {
assertThat(e).isInstanceOf(ExecutionException.class);
assertThat(e).hasCauseThat().isInstanceOf(PermissionDeniedException.class);
return true;
});
}

// --- appendEvent Tests ---
/** Tests that appendEvent persists the event and updates the session's updateTime. */
@Test
Expand Down Expand Up @@ -978,4 +1094,12 @@ void deleteSession_appNameMismatch_doesNotDelete() {
verify(mockSessionDocRef, never()).delete();
verify(mockWriteBatch, never()).commit();
}

private static AlreadyExistsException alreadyExists() {
return new AlreadyExistsException(
"Document already exists",
/* cause= */ null,
GrpcStatusCode.of(Status.Code.ALREADY_EXISTS),
/* retryable= */ false);
}
}
Loading
Loading