/
githubmirror
/
galahad
Обзор
Документация
Войти
/
githubmirror
/
galahad
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/java.net.http/share/classes/jdk/internal/net/http/PullPublisher.java
166 строк
6 KB
Volkan Yazici
8367067: Improve exception handling in HttpRequest.BodyPublishers
19 сен 2025, 15:07
19 сен 2025, 15:07
87d5042
Код
Авторство
О чём код?
/* * Copyright (c) 2016, 2025, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it * under the terms of the GNU General Public License version 2 only, as * published by the Free Software Foundation. Oracle designates this * particular file as subject to the "Classpath" exception as provided * by Oracle in the LICENSE file that accompanied this code. * * This code is distributed in the hope that it will be useful, but WITHOUT * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License * version 2 for more details (a copy is included in the LICENSE file that * accompanied this code). * * You should have received a copy of the GNU General Public License version * 2 along with this work; if not, write to the Free Software Foundation, * Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA. * * Please contact Oracle, 500 Oracle Parkway, Redwood Shores, CA 94065 USA * or visit www.oracle.com if you need additional information or have any * questions. */ package jdk.internal.net.http; import java.util.concurrent.Flow; import jdk.internal.net.http.common.Demand; import jdk.internal.net.http.common.SequentialScheduler; /** * A {@linkplain Flow.Publisher publisher} that publishes items obtained from the given {@link CheckedIterable}. * Each new subscription gets a new {@link CheckedIterator}. */ class PullPublisher<T> implements Flow.Publisher<T> { // Only one of `iterable` or `throwable` should be null, and the other non-null. throwable is // non-null when an error has been encountered, by the creator of // PullPublisher, while subscribing the subscriber, but before subscribe has // completed. private final CheckedIterable<T> iterable; private final Throwable throwable; PullPublisher(CheckedIterable<T> iterable, Throwable throwable) { if ((iterable == null) == (throwable == null)) { String message = String.format( "only one of `iterable` or `throwable` should be null, and the other non-null, but %s are null", throwable == null ? "both" : "none"); throw new IllegalArgumentException(message); } this.iterable = iterable; this.throwable = throwable; } PullPublisher(CheckedIterable<T> iterable) { this(iterable, null); } @Override public void subscribe(Flow.Subscriber<? super T> subscriber) { Throwable failure = throwable; CheckedIterator<T> iterator = null; if (failure == null) { try { iterator = iterable.iterator(); } catch (Exception exception) { failure = exception; } } Subscription sub = failure != null ? new Subscription(subscriber, null, failure) : new Subscription(subscriber, iterator, null); subscriber.onSubscribe(sub); if (failure != null) { sub.pullScheduler.runOrSchedule(); } } private class Subscription implements Flow.Subscription { private final Flow.Subscriber<? super T> subscriber; private final CheckedIterator<T> iter; private volatile boolean completed; private volatile boolean cancelled; private volatile Throwable error; final SequentialScheduler pullScheduler = new SequentialScheduler(new PullTask()); private final Demand demand = new Demand(); Subscription(Flow.Subscriber<? super T> subscriber, CheckedIterator<T> iter, Throwable throwable) { this.subscriber = subscriber; this.iter = iter; this.error = throwable; } final class PullTask extends SequentialScheduler.CompleteRestartableTask { @Override protected void run() { if (completed || cancelled) { return; } Throwable t = error; if (t != null) { completed = true; pullScheduler.stop(); subscriber.onError(t); return; } while (demand.tryDecrement() && !cancelled) { T next; try { if (!iter.hasNext()) { break; } next = iter.next(); } catch (Throwable t1) { completed = true; pullScheduler.stop(); subscriber.onError(t1); return; } subscriber.onNext(next); } boolean hasNext; try { hasNext = iter.hasNext(); } catch (Exception e) { completed = true; pullScheduler.stop(); subscriber.onError(e); return; } if (!hasNext && !cancelled) { completed = true; pullScheduler.stop(); subscriber.onComplete(); } } } @Override public void request(long n) { if (cancelled) return; // no-op if (n <= 0) { error = new IllegalArgumentException("non-positive subscription request: " + n); } else { demand.increase(n); } pullScheduler.runOrSchedule(); } @Override public void cancel() { cancelled = true; } } }