Repository navigation
Feat/add teams channel plugin #96
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,13 @@ | ||
| plugins { | ||
| id 'java-library' | ||
| } | ||
|
|
||
| dependencies { | ||
| implementation project(':base') | ||
| implementation 'org.springframework.boot:spring-boot-starter' | ||
| implementation 'org.springframework.boot:spring-boot-starter-webmvc' | ||
| implementation 'com.nimbusds:nimbus-jose-jwt:10.9.1' | ||
|
|
||
| testImplementation 'org.springframework.boot:spring-boot-starter-test' | ||
| testImplementation 'org.springframework.boot:spring-boot-starter-json' | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,88 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import com.nimbusds.jose.JWSAlgorithm; | ||
| import com.nimbusds.jose.jwk.source.JWKSource; | ||
| import com.nimbusds.jose.jwk.source.JWKSourceBuilder; | ||
| import com.nimbusds.jose.proc.JWSVerificationKeySelector; | ||
| import com.nimbusds.jose.proc.SecurityContext; | ||
| import com.nimbusds.jwt.JWTClaimsSet; | ||
| import com.nimbusds.jwt.SignedJWT; | ||
| import com.nimbusds.jwt.proc.DefaultJWTProcessor; | ||
| import com.nimbusds.jwt.proc.JWTProcessor; | ||
| import org.slf4j.Logger; | ||
| import org.slf4j.LoggerFactory; | ||
|
|
||
| import java.net.MalformedURLException; | ||
| import java.net.URI; | ||
| import java.time.Duration; | ||
| import java.util.Date; | ||
|
|
||
| /** | ||
| * Validates the Bot Framework Connector JWT sent on every incoming Teams webhook call | ||
| * (the {@code Authorization: Bearer <token>} header), per | ||
| * https://learn.microsoft.com/azure/bot-service/rest-api/bot-framework-rest-connector-authentication | ||
| * <p> | ||
| * A valid token must: (1) be signed by a key currently published in Bot Framework's JWKS, | ||
| * (2) have issuer {@code https://api.botframework.com}, (3) have this bot's App ID as | ||
| * audience, and (4) not be expired. The JWKS itself is fetched once and cached/refreshed | ||
| * automatically by the underlying {@link JWKSource}. | ||
| */ | ||
| class BotFrameworkJwtValidator { | ||
|
|
||
| private static final Logger LOGGER = LoggerFactory.getLogger(BotFrameworkJwtValidator.class); | ||
| private static final String EXPECTED_ISSUER = "https://api.botframework.com"; | ||
| private static final String JWKS_URL = "https://login.botframework.com/v1/.well-known/keys"; | ||
|
|
||
| private final TeamsProperties properties; | ||
| private final JWTProcessor<SecurityContext> jwtProcessor; | ||
|
|
||
| BotFrameworkJwtValidator(TeamsProperties properties) { | ||
| this(properties, defaultProcessor()); | ||
| } | ||
|
|
||
| /** Test seam: inject a stub/mocked processor instead of hitting the real JWKS endpoint. */ | ||
| BotFrameworkJwtValidator(TeamsProperties properties, JWTProcessor<SecurityContext> jwtProcessor) { | ||
| this.properties = properties; | ||
| this.jwtProcessor = jwtProcessor; | ||
| } | ||
|
|
||
| private static JWTProcessor<SecurityContext> defaultProcessor() { | ||
| try { | ||
| JWKSource<SecurityContext> jwkSource = JWKSourceBuilder | ||
| .create(URI.create(JWKS_URL).toURL()) | ||
| .cache(Duration.ofMinutes(30).toMillis(), Duration.ofMinutes(1).toMillis()) | ||
| .build(); | ||
| DefaultJWTProcessor<SecurityContext> processor = new DefaultJWTProcessor<>(); | ||
| processor.setJWSKeySelector(new JWSVerificationKeySelector<>(JWSAlgorithm.RS256, jwkSource)); | ||
| return processor; | ||
| } catch (MalformedURLException e) { | ||
| throw new IllegalStateException("Invalid Bot Framework JWKS URL: " + JWKS_URL, e); | ||
| } | ||
| } | ||
|
|
||
| /** @return true if {@code bearerToken} is a currently valid Bot Framework Connector token for this bot */ | ||
| boolean isValid(String bearerToken) { | ||
| try { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The Bot Framework spec also requires the token’s serviceUrl claim to match the activity’s serviceUrl. That check is important here, because we send the bot’s outbound access token to activity.serviceUrl(). Please pass the activity’s serviceUrl into isValid(...) and compare it to the claim. Separately, the issuer, audience and expiry checks can be handed to Nimbus, which handles clock skew: |
||
| SignedJWT jwt = SignedJWT.parse(bearerToken); | ||
| JWTClaimsSet claims = jwtProcessor.process(jwt, null); // throws if signature/algorithm is invalid | ||
|
|
||
| if (!EXPECTED_ISSUER.equals(claims.getIssuer())) { | ||
| LOGGER.warn("Rejected Teams webhook JWT with unexpected issuer '{}'", claims.getIssuer()); | ||
| return false; | ||
| } | ||
| if (properties.getAppId() == null || !claims.getAudience().contains(properties.getAppId())) { | ||
| LOGGER.warn("Rejected Teams webhook JWT with unexpected audience {}", claims.getAudience()); | ||
| return false; | ||
| } | ||
| Date expiration = claims.getExpirationTime(); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This exp check duplicates what DefaultJWTProcessor already does, and it has no clock-skew tolerance. If the claims verifier above is added, this check (and the issuer and audience checks at lines 69–76) can be removed. |
||
| if (expiration == null || expiration.before(new Date())) { | ||
| LOGGER.warn("Rejected expired Teams webhook JWT"); | ||
| return false; | ||
| } | ||
| return true; | ||
| } catch (Exception e) { | ||
| LOGGER.warn("Rejected Teams webhook JWT: {}", e.getMessage()); | ||
| return false; | ||
| } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,79 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import tools.jackson.databind.json.JsonMapper; | ||
| import org.slf4j.Logger; | ||
| import org.slf4j.LoggerFactory; | ||
|
|
||
| import java.net.URI; | ||
| import java.net.http.HttpClient; | ||
| import java.net.http.HttpRequest; | ||
| import java.net.http.HttpResponse; | ||
| import java.time.Instant; | ||
| import java.util.concurrent.atomic.AtomicReference; | ||
|
|
||
| /** | ||
| * Fetches and caches an OAuth2 client-credentials token for the Bot Framework | ||
| * Connector API, using the Azure AD v2 endpoint (login.microsoftonline.com). | ||
| */ | ||
| class BotFrameworkTokenProvider { | ||
|
|
||
| private static final Logger LOGGER = LoggerFactory.getLogger(BotFrameworkTokenProvider.class); | ||
| private static final String SCOPE = "https://api.botframework.com/.default"; | ||
| private static final JsonMapper JSON = JsonMapper.builder().build(); | ||
|
|
||
| private final TeamsProperties properties; | ||
| private final HttpClient httpClient; | ||
| private final AtomicReference<CachedToken> cached = new AtomicReference<>(); | ||
|
|
||
| BotFrameworkTokenProvider(TeamsProperties properties, HttpClient httpClient) { | ||
| this.properties = properties; | ||
| this.httpClient = httpClient; | ||
| } | ||
|
|
||
| /** @return a valid bearer token, refreshing it if expired or missing */ | ||
| synchronized String getToken() { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. synchronized together with AtomicReference (line 26) is redundant. Since the method is already synchronized, a plain field is enough. |
||
| CachedToken token = cached.get(); | ||
| if (token != null && Instant.now().isBefore(token.expiresAt())) { | ||
| return token.value(); | ||
| } | ||
| CachedToken fresh = fetchToken(); | ||
| cached.set(fresh); | ||
| return fresh.value(); | ||
| } | ||
|
|
||
| private CachedToken fetchToken() { | ||
| String tokenUrl = "https://login.microsoftonline.com/" + properties.getTenantId() | ||
| + "/oauth2/v2.0/token"; | ||
| String body = "grant_type=client_credentials" | ||
| + "&client_id=" + properties.getAppId() | ||
| + "&client_secret=" + properties.getAppSecret() | ||
|
Comment on lines
+48
to
+49
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The form body isn’t URL-encoded. If the secret contains reserved characters, the token request fails or is misparsed: |
||
| + "&scope=" + SCOPE; | ||
|
|
||
| HttpRequest request = HttpRequest.newBuilder(URI.create(tokenUrl)) | ||
| .header("Content-Type", "application/x-www-form-urlencoded") | ||
| .POST(HttpRequest.BodyPublishers.ofString(body)) | ||
| .build(); | ||
| try { | ||
| HttpResponse<String> response = httpClient.send(request, HttpResponse.BodyHandlers.ofString()); | ||
| if (response.statusCode() != 200) { | ||
| throw new IllegalStateException( | ||
| "Bot Framework token request failed with status " + response.statusCode() | ||
| + ": " + response.body()); | ||
| } | ||
| TokenResponse parsed = JSON.readValue(response.body(), TokenResponse.class); | ||
| return new CachedToken(parsed.accessToken(), | ||
| Instant.now().plusSeconds(Math.max(0, parsed.expiresIn() - 60))); | ||
| } catch (Exception e) { | ||
| LOGGER.error("Failed to fetch Bot Framework access token", e); | ||
| throw new IllegalStateException("Could not obtain Bot Framework access token", e); | ||
| } | ||
| } | ||
|
|
||
| private record CachedToken(String value, Instant expiresAt) { | ||
| } | ||
|
|
||
| private record TokenResponse( | ||
| @com.fasterxml.jackson.annotation.JsonProperty("access_token") String accessToken, | ||
| @com.fasterxml.jackson.annotation.JsonProperty("expires_in") int expiresIn) { | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,23 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import com.fasterxml.jackson.annotation.JsonIgnoreProperties; | ||
| import com.fasterxml.jackson.annotation.JsonProperty; | ||
|
|
||
| @JsonIgnoreProperties(ignoreUnknown = true) | ||
| public record TeamsActivityPayload( | ||
| @JsonProperty("type") String type, | ||
| @JsonProperty("id") String id, | ||
| @JsonProperty("serviceUrl") String serviceUrl, | ||
| @JsonProperty("channelId") String channelId, | ||
| @JsonProperty("text") String text, | ||
| @JsonProperty("from") From from, | ||
| @JsonProperty("conversation") Conversation conversation) { | ||
|
|
||
| @JsonIgnoreProperties(ignoreUnknown = true) | ||
| public record From(@JsonProperty("id") String id, @JsonProperty("name") String name) { | ||
| } | ||
|
|
||
| @JsonIgnoreProperties(ignoreUnknown = true) | ||
| public record Conversation(@JsonProperty("id") String id) { | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,79 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import ai.javaclaw.channels.Channel; | ||
| import ai.javaclaw.channels.ChannelRegistry; | ||
| import org.slf4j.Logger; | ||
| import org.slf4j.LoggerFactory; | ||
| import tools.jackson.databind.json.JsonMapper; | ||
|
|
||
| import java.net.URI; | ||
| import java.net.http.HttpClient; | ||
| import java.net.http.HttpRequest; | ||
| import java.net.http.HttpResponse; | ||
| import java.util.Map; | ||
| import java.util.concurrent.atomic.AtomicReference; | ||
|
|
||
| public class TeamsChannel implements Channel { | ||
|
|
||
| private static final Logger LOGGER = LoggerFactory.getLogger(TeamsChannel.class); | ||
| private static final JsonMapper JSON = JsonMapper.builder().build(); | ||
|
|
||
| static final String CHANNEL_ID = "teams"; | ||
|
|
||
| private final BotFrameworkTokenProvider tokenProvider; | ||
| private final HttpClient httpClient; | ||
| private final AtomicReference<TeamsConversationReference> conversationReference = | ||
| new AtomicReference<>(); | ||
|
|
||
| public TeamsChannel(BotFrameworkTokenProvider tokenProvider, ChannelRegistry channelRegistry) { | ||
| this(tokenProvider, HttpClient.newHttpClient(), channelRegistry); | ||
| } | ||
|
|
||
| TeamsChannel(BotFrameworkTokenProvider tokenProvider, HttpClient httpClient, ChannelRegistry channelRegistry) { | ||
| this.tokenProvider = tokenProvider; | ||
| this.httpClient = httpClient; | ||
| channelRegistry.registerChannel(this); | ||
| LOGGER.info("Started Microsoft Teams channel"); | ||
| } | ||
|
|
||
| @Override | ||
| public String getName() { | ||
| return CHANNEL_ID; | ||
| } | ||
|
|
||
| /** Called by the webhook controller whenever a message arrives from the allowed user. */ | ||
| void updateConversationReference(TeamsConversationReference reference) { | ||
| conversationReference.set(reference); | ||
| } | ||
|
|
||
| @Override | ||
| public void sendMessage(String message) { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Related to the comment above: please add sendMessage(TeamsConversationReference ref, String message) and have sendMessage(String) delegate to it with the last known reference. |
||
| if (message == null || message.isBlank()) { | ||
| return; | ||
| } | ||
| TeamsConversationReference reference = conversationReference.get(); | ||
| if (reference == null) { | ||
| LOGGER.warn("No known Teams conversation yet, cannot send message '{}'", message); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: these log the full message content at warn/error level. Consider leaving it out or truncating it. |
||
| return; | ||
| } | ||
|
|
||
| String url = reference.serviceUrl() | ||
| + (reference.serviceUrl().endsWith("/") ? "" : "/") | ||
| + "v3/conversations/" + reference.conversationId() + "/activities"; | ||
|
|
||
| try { | ||
| String body = JSON.writeValueAsString(Map.of("type", "message", "text", message)); | ||
| HttpRequest request = HttpRequest.newBuilder(URI.create(url)) | ||
| .header("Authorization", "Bearer " + tokenProvider.getToken()) | ||
| .header("Content-Type", "application/json") | ||
| .POST(HttpRequest.BodyPublishers.ofString(body)) | ||
| .build(); | ||
| HttpResponse<String> response = httpClient.send(request, HttpResponse.BodyHandlers.ofString()); | ||
| if (response.statusCode() >= 300) { | ||
| LOGGER.error("Teams send failed with status {}: {}", response.statusCode(), response.body()); | ||
| } | ||
| } catch (Exception e) { | ||
| LOGGER.error("Failed to send Teams message '{}'", message, e); | ||
| } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,57 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import ai.javaclaw.channels.ChannelRegistry; | ||
| import org.springframework.boot.autoconfigure.AutoConfiguration; | ||
| import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; | ||
| import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; | ||
| import org.springframework.boot.context.properties.EnableConfigurationProperties; | ||
| import org.springframework.context.annotation.Bean; | ||
| import ai.javaclaw.agent.Agent; | ||
|
|
||
| import java.net.http.HttpClient; | ||
|
|
||
| @AutoConfiguration | ||
| @EnableConfigurationProperties(TeamsProperties.class) | ||
| public class TeamsChannelAutoConfiguration { | ||
|
|
||
| @Bean | ||
| @ConditionalOnMissingBean | ||
| public BotFrameworkTokenProvider botFrameworkTokenProvider(TeamsProperties properties) { | ||
| return new BotFrameworkTokenProvider(properties, HttpClient.newHttpClient()); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. There’s no connect timeout here, and no request timeouts on the send calls (BotFrameworkTokenProvider.java:57, TeamsChannel.java:71). With a single worker thread, one hung call to Azure blocks every Teams message after it. Suggest: plus .timeout(Duration.ofSeconds(30)) on each HttpRequest. The same applies to TeamsChannel.java:29. |
||
| } | ||
|
|
||
| @Bean | ||
| @ConditionalOnMissingBean | ||
| public BotFrameworkJwtValidator botFrameworkJwtValidator(TeamsProperties properties) { | ||
| return new BotFrameworkJwtValidator(properties); | ||
|
Comment on lines
+20
to
+26
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. These two beans are created even when Teams is disabled. Please add the same @ConditionalOnProperty(prefix = "agent.channels.teams", name = "enabled", havingValue = "true") used on the other beans. |
||
| } | ||
|
|
||
| @Bean | ||
| @ConditionalOnMissingBean | ||
| @ConditionalOnProperty(prefix = "agent.channels.teams", name = "enabled", havingValue = "true") | ||
| public TeamsChannel teamsChannel(BotFrameworkTokenProvider tokenProvider, ChannelRegistry channelRegistry) { | ||
| return new TeamsChannel(tokenProvider, channelRegistry); | ||
| } | ||
|
|
||
| @Bean | ||
| @ConditionalOnProperty( | ||
| prefix = "agent.channels.teams", | ||
| name = "enabled", | ||
| havingValue = "true" | ||
| ) | ||
| public TeamsWebhookController teamsWebhookController( | ||
| TeamsProperties properties, | ||
| ChannelRegistry channelRegistry, | ||
| Agent agent, | ||
| TeamsChannel channel, | ||
| BotFrameworkJwtValidator jwtValidator) { | ||
|
|
||
| return new TeamsWebhookController( | ||
| properties, | ||
| channelRegistry, | ||
| agent, | ||
| channel, | ||
| jwtValidator | ||
| ); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. TeamsWebhookController is already a @RestController in ai.javaclaw.*, so component scanning picks it up. This @bean silently replaces the scanned definition. It works, but it’s redundant. WhatsApp relies on scanning only, so I’d remove this method. |
||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,17 @@ | ||
| package ai.javaclaw.channels.teams; | ||
|
|
||
| import ai.javaclaw.channels.ChannelMessageReceivedEvent; | ||
|
|
||
| public class TeamsChannelMessageReceivedEvent extends ChannelMessageReceivedEvent { | ||
|
|
||
| private final String conversationId; | ||
|
|
||
| public TeamsChannelMessageReceivedEvent(String channel, String message, String conversationId) { | ||
| super(channel, message); | ||
| this.conversationId = conversationId; | ||
| } | ||
|
|
||
| public String getConversationId() { | ||
| return conversationId; | ||
| } | ||
| } |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Nit: Spring Boot already manages the nimbus-jose-jwt version, so :10.9.1 can be dropped.