11package io .github .easy4j .opencode .api .sse ;
22
33import java .util .Objects ;
4+ import java .util .concurrent .CancellationException ;
5+ import java .util .concurrent .CompletableFuture ;
46import java .util .concurrent .atomic .AtomicBoolean ;
57
68/**
@@ -20,13 +22,35 @@ public final class SseSubscription implements AutoCloseable {
2022 */
2123 private final Runnable cancellation ;
2224
25+ /** Completes when the SSE transport is connected and ready to receive events. */
26+ private final CompletableFuture <Void > ready ;
27+
28+ /** Completes with the terminal transport failure when one occurs. */
29+ private final CompletableFuture <Throwable > failure ;
30+
2331 /**
2432 * 创建 sse subscription 实例,并按传入依赖确定资源所有权。
2533 *
2634 * @param cancellation 取消信号;为 {@code null} 时不可由外部取消
2735 */
2836 public SseSubscription (Runnable cancellation ) {
37+ this (cancellation , new CompletableFuture <>(), new CompletableFuture <>());
38+ }
39+
40+ public SseSubscription (Runnable cancellation ,
41+ CompletableFuture <Void > ready ,
42+ CompletableFuture <Throwable > failure ) {
2943 this .cancellation = Objects .requireNonNull (cancellation , "cancellation" );
44+ this .ready = Objects .requireNonNull (ready , "ready" );
45+ this .failure = Objects .requireNonNull (failure , "failure" );
46+ }
47+
48+ public CompletableFuture <Void > getReady () {
49+ return ready ;
50+ }
51+
52+ public CompletableFuture <Throwable > getFailure () {
53+ return failure ;
3054 }
3155
3256 /**
@@ -38,6 +62,9 @@ public boolean cancel() {
3862 if (!active .compareAndSet (true , false )) {
3963 return false ;
4064 }
65+ if (!ready .isDone ()) {
66+ ready .completeExceptionally (new CancellationException ("SSE subscription cancelled before readiness" ));
67+ }
4168 cancellation .run ();
4269 return true ;
4370 }
0 commit comments