poolster-plugin-java 0.5.0-alpha.1

Native Java SDK generator for Poolster
Documentation
package __PACKAGE__;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.net.URI;
import java.net.URLEncoder;
import java.net.http.*;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.Flow;
import java.util.concurrent.atomic.AtomicReference;
/** Client-credentials only. Separate issuer transport must disable redirects. */
public final class OAuthClientCredentials {
    private final HttpClient issuer;
    private final URI tokenUrl;
    private final String basic, scope;
    private final Duration leeway, timeout;
    private final Semaphore refresh = new Semaphore(1);
    private final AtomicReference<Token> cached = new AtomicReference<>();
    private record Token(String value, long expires) {}
    public OAuthClientCredentials(HttpClient issuer,String tokenUrl,String clientId,String clientSecret,List<String> scopes,Duration refreshLeeway,Duration timeout) {
        this.issuer=Objects.requireNonNull(issuer);this.tokenUrl=URI.create(tokenUrl);
        if(!Set.of("https","http").contains(this.tokenUrl.getScheme()) || this.tokenUrl.getHost()==null || this.tokenUrl.getUserInfo()!=null || this.tokenUrl.getFragment()!=null || issuer.followRedirects()!=HttpClient.Redirect.NEVER || clientId==null || clientId.isEmpty() || clientSecret==null || clientSecret.isEmpty())throw new IllegalArgumentException("Invalid OAuth configuration");
        this.leeway=refreshLeeway==null?Duration.ofSeconds(30):refreshLeeway;this.timeout=timeout==null?Duration.ofSeconds(30):timeout;
        if(this.leeway.isNegative()||this.timeout.isNegative()||this.timeout.isZero())throw new IllegalArgumentException("Invalid OAuth configuration");
        this.basic=Base64.getEncoder().encodeToString((URLEncoder.encode(clientId,StandardCharsets.UTF_8)+":"+URLEncoder.encode(clientSecret,StandardCharsets.UTF_8)).getBytes(StandardCharsets.US_ASCII));
        this.scope=String.join(" ",scopes==null?List.of():List.copyOf(scopes));
    }
    public OAuthClientCredentials(HttpClient issuer,String tokenUrl,String clientId,String clientSecret) {this(issuer,tokenUrl,clientId,clientSecret,List.of(),null,null);}
    public void invalidate(String rejected) {cached.updateAndGet(value->value!=null&&value.value().equals(rejected)?null:value);}
    public String token() throws IOException,InterruptedException {
        if(Thread.currentThread().isInterrupted())throw new InterruptedException();
        Token value=cached.get();if(value!=null&&System.nanoTime()<value.expires())return value.value();
        refresh.acquire();
        try {
            value=cached.get();if(value!=null&&System.nanoTime()<value.expires())return value.value();
            String body="grant_type=client_credentials"+(scope.isEmpty()?"":"&scope="+URLEncoder.encode(scope,StandardCharsets.UTF_8));
            var request=HttpRequest.newBuilder(tokenUrl).timeout(timeout).header("Authorization","Basic "+basic).header("Content-Type","application/x-www-form-urlencoded").POST(HttpRequest.BodyPublishers.ofString(body)).build();
            var response=issuer.send(request,info->new LimitedBody());
            if(response.statusCode()<200||response.statusCode()>=300)throw new IOException();
            JsonNode json=new ObjectMapper().readTree(response.body());
            if(json==null||!json.path("access_token").isTextual()||!json.path("token_type").isTextual()||!json.path("token_type").asText().equalsIgnoreCase("Bearer")||!json.path("expires_in").isNumber())throw new IOException();
            String token=json.path("access_token").asText();double seconds=json.path("expires_in").asDouble();
            if(token.isEmpty()||!token.matches("[A-Za-z0-9._~+/-]+=*")||!Double.isFinite(seconds)||seconds<0)throw new IOException();
            if(Thread.currentThread().isInterrupted())throw new InterruptedException();
            double duration=Math.max(0,seconds-leeway.toMillis()/1000.0);if(duration>0) {
                // Bound issuer lifetimes to one year, preserving monotonic expiry arithmetic.
                long nanos=(long)(Math.min(duration,365.0*24*3600)*1_000_000_000);cached.set(new Token(token,System.nanoTime()+nanos));
            }
            return token;
        } catch(InterruptedException e){throw e;}catch(Exception e){throw new IOException("OAuth token acquisition failed");}finally{refresh.release();}
    }
    @Override public String toString(){return "OAuthClientCredentials";}
    private static final class LimitedBody implements HttpResponse.BodySubscriber<byte[]> {
        private final CompletableFuture<byte[]> result=new CompletableFuture<>();private final java.io.ByteArrayOutputStream bytes=new java.io.ByteArrayOutputStream();private Flow.Subscription subscription;
        public CompletionStage<byte[]> getBody(){return result;}
        public void onSubscribe(Flow.Subscription value){subscription=value;value.request(1);}
        public void onNext(List<ByteBuffer> values){for(var value:values){if(bytes.size()+value.remaining()>65536){subscription.cancel();result.completeExceptionally(new IOException("OAuth token acquisition failed"));return;}byte[] part=new byte[value.remaining()];value.get(part);bytes.writeBytes(part);}subscription.request(1);}
        public void onError(Throwable error){result.completeExceptionally(new IOException("OAuth token acquisition failed"));}
        public void onComplete(){result.complete(bytes.toByteArray());}
    }
}