From 7975d60abadeab29115dc7f0952dbfafb1336294 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:01:32 -0400 Subject: [PATCH 1/9] fix(avio): refresh direct-play credentials on every range request Direct-play sources sent the headers they were opened with for the whole session. LoadOptions.httpRequestAuthorization only reached native HLS, so when a host rotated its access token every later range request was refused, and the host had to reload the player at the current position to send the new Authorization header. The byte-range reader now accepts the same HTTPRequestAuthorization and asks it for the source URL before each request it builds: pump ranges, reconnects, seeks, detour blocks, size probes, the tail prefetch and the streaming GET. The answer replaces the static headers and then follows RedirectHeaderPolicy, so credentials never reach a cross-origin redirect target. A 401 asks the provider once with the headers that request carried, and a changed Authorization retries once at the same byte offset. Unchanged credentials, a second 401, or a provider that throws or exceeds its 10 s bound fail the read without running the reconnect ladder, and fail an open before anything is sent. Closing the reader cancels pending waits. The provider reaches the playback demuxer, its reopens and reloads, and the embedded-subtitle side readers. Static headers stay the default. Co-Authored-By: Claude Opus 5.5 (1M context) --- .github/workflows/ci.yml | 2 +- CHANGELOG.md | 2 + CONTRIBUTING.md | 6 +- .../AetherEngine/AetherEngine+Loading.swift | 18 +- .../AetherEngine/AetherEngine+Subtitles.swift | 15 +- Sources/AetherEngine/AetherEngine.swift | 6 +- Sources/AetherEngine/Demuxer/AVIOReader.swift | 264 ++++++++++++++++-- Sources/AetherEngine/Demuxer/Demuxer.swift | 10 +- .../AetherEngine/Network/HLSOriginRelay.swift | 12 +- .../Network/HTTPRequestAuthorization.swift | 80 +++++- Sources/AetherEngine/PlayerState.swift | 8 +- .../Video/HLSVideoEngine+LiveReopen.swift | 3 +- .../AetherEngine/Video/HLSVideoEngine.swift | 8 +- .../Issue255BodyReserveTests.swift | 5 +- ...reshableDirectPlayAuthorizationTests.swift | 213 ++++++++++++++ docs/api.md | 16 +- 16 files changed, 598 insertions(+), 70 deletions(-) create mode 100644 Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index cd6d81930..48d7c25ad 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -17,7 +17,7 @@ jobs: runs-on: macos-15 timeout-minutes: 30 env: - AUTHORIZATION_TEST_SUITES: 'RefreshableHLSAuthorizationTests|RefreshableSubtitleAuthorizationTests|LiveTrustEvaluatorTests' + AUTHORIZATION_TEST_SUITES: 'RefreshableHLSAuthorizationTests|RefreshableSubtitleAuthorizationTests|RefreshableDirectPlayAuthorizationTests|LiveTrustEvaluatorTests' steps: - uses: actions/checkout@v4 diff --git a/CHANGELOG.md b/CHANGELOG.md index c91fd1c81..ba3830434 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,6 +23,8 @@ the public-API contract. - `ExternalSubtitleTrack.httpRequestAuthorization` supplies refreshable headers for primary/secondary sidecars and native subtitle stores without changing registered track IDs or rendition mappings. Authorized container decoding retains AVIO streaming and range access. - `HTTPRequestAuthorization.data(from:maximumBytes:)` fetches raw auxiliary resources such as font bundles with a caller-supplied byte limit and a whole-transfer deadline, reusing the relay's redirect, authorization, retry, cancellation and TLS policy. +- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest and remote disc images keep static headers. + - `LoadOptions.httpRequestAuthorization` accepts an async `HTTPRequestAuthorization` resolver for native HLS. The engine resolves headers before requests and redirects, and retries a rejected request once when the bearer changes, preserving the active player item across token rotation. ### Fixed diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index dacf79664..a900dcc1e 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -17,8 +17,8 @@ swift build swift test ``` -CI runs `RefreshableHLSAuthorizationTests`, `RefreshableSubtitleAuthorizationTests` and -`LiveTrustEvaluatorTests` in a separate process because their short authorization deadlines require +CI runs `RefreshableHLSAuthorizationTests`, `RefreshableSubtitleAuthorizationTests`, +`RefreshableDirectPlayAuthorizationTests` and `LiveTrustEvaluatorTests` in a separate process because their short authorization deadlines require responsive async resolvers. Blocking work elsewhere in the suite can delay those resolvers on smaller runners. `LiveTrustEvaluatorTests` is the serialized parent of every live suite that sets the process-global `EngineTLS.serverTrustEvaluator`. Nest any new suite that sets it there; the @@ -28,7 +28,7 @@ The two commands below cover the entire test suite, keeping the existing deadlin and parallel execution within each group: ```bash -AUTHORIZATION_TEST_SUITES='RefreshableHLSAuthorizationTests|RefreshableSubtitleAuthorizationTests|LiveTrustEvaluatorTests' +AUTHORIZATION_TEST_SUITES='RefreshableHLSAuthorizationTests|RefreshableSubtitleAuthorizationTests|RefreshableDirectPlayAuthorizationTests|LiveTrustEvaluatorTests' swift test --skip "$AUTHORIZATION_TEST_SUITES" swift test --skip-build --filter "$AUTHORIZATION_TEST_SUITES" ``` diff --git a/Sources/AetherEngine/AetherEngine+Loading.swift b/Sources/AetherEngine/AetherEngine+Loading.swift index f79aae82e..d035904cc 100644 --- a/Sources/AetherEngine/AetherEngine+Loading.swift +++ b/Sources/AetherEngine/AetherEngine+Loading.swift @@ -692,6 +692,7 @@ extension AetherEngine { func loadNative( url: URL, sourceHTTPHeaders: [String: String] = [:], + sourceHTTPAuthorization: HTTPRequestAuthorization? = nil, startPosition: Double?, audioSourceStreamIndex: Int32? = nil, keepDvh1TagWithoutDV: Bool = false, @@ -760,6 +761,7 @@ extension AetherEngine { let session = HLSVideoEngine( url: url, sourceHTTPHeaders: sourceHTTPHeaders, + sourceHTTPAuthorization: sourceHTTPAuthorization, dvModeAvailable: sessionDisplayCaps.supportsDolbyVision, displaySupportsHDR: sessionDisplayCaps.supportsHDR, keepDvh1TagWithoutDV: keepDvh1TagWithoutDV, @@ -1725,6 +1727,7 @@ extension AetherEngine { func loadSoftware( url: URL, sourceHTTPHeaders: [String: String] = [:], + sourceHTTPAuthorization: HTTPRequestAuthorization? = nil, startPosition: Double?, audioSourceStreamIndex: Int32?, isLive: Bool = false, @@ -1876,13 +1879,13 @@ extension AetherEngine { if loadGeneration == generation { recordStartupCheckpoint(.sessionConstructed) } // #361 let forwardBufferSegments = loadedOptions.forwardBufferSegments try await Task.detached(priority: .userInitiated) { - [host, preopenedDemuxer, url, sourceHTTPHeaders, isLive, dvrWindowSeconds, probesize, maxAnalyzeDuration, sequentialOrigin, heldSourceConnection, declaredDuration, networkPhaseSink] in + [host, preopenedDemuxer, url, sourceHTTPHeaders, sourceHTTPAuthorization, isLive, dvrWindowSeconds, probesize, maxAnalyzeDuration, sequentialOrigin, heldSourceConnection, declaredDuration, networkPhaseSink] in let dem: Demuxer if let pre = preopenedDemuxer { dem = pre } else { dem = Demuxer() - try dem.open(url: url, extraHeaders: sourceHTTPHeaders, profile: .playback.withProbeBudget(probesize: probesize, maxAnalyzeDuration: maxAnalyzeDuration).withSequentialOrigin(sequentialOrigin, declaredDuration: declaredDuration).withHeldSourceConnection(heldSourceConnection), isLive: isLive) + try dem.open(url: url, extraHeaders: sourceHTTPHeaders, requestAuthorization: sourceHTTPAuthorization, profile: .playback.withProbeBudget(probesize: probesize, maxAnalyzeDuration: maxAnalyzeDuration).withSequentialOrigin(sequentialOrigin, declaredDuration: declaredDuration).withHeldSourceConnection(heldSourceConnection), isLive: isLive) } dem.onNetworkPhaseChanged = networkPhaseSink try await host.load( @@ -1905,6 +1908,7 @@ extension AetherEngine { func loadAudio( url: URL, sourceHTTPHeaders: [String: String] = [:], + sourceHTTPAuthorization: HTTPRequestAuthorization? = nil, startPosition: Double?, audioSourceStreamIndex: Int32?, preopenedDemuxer: Demuxer?, @@ -1953,13 +1957,13 @@ extension AetherEngine { } if loadGeneration == generation { recordStartupCheckpoint(.sessionConstructed) } // #361 try await Task.detached(priority: .userInitiated) { - [host, preopenedDemuxer, url, sourceHTTPHeaders, probesize, maxAnalyzeDuration, sequentialOrigin, heldSourceConnection, declaredDuration, networkPhaseSink] in + [host, preopenedDemuxer, url, sourceHTTPHeaders, sourceHTTPAuthorization, probesize, maxAnalyzeDuration, sequentialOrigin, heldSourceConnection, declaredDuration, networkPhaseSink] in let dem: Demuxer if let pre = preopenedDemuxer { dem = pre } else { dem = Demuxer() - try dem.open(url: url, extraHeaders: sourceHTTPHeaders, profile: .playback.withProbeBudget(probesize: probesize, maxAnalyzeDuration: maxAnalyzeDuration).withSequentialOrigin(sequentialOrigin, declaredDuration: declaredDuration).withHeldSourceConnection(heldSourceConnection)) + try dem.open(url: url, extraHeaders: sourceHTTPHeaders, requestAuthorization: sourceHTTPAuthorization, profile: .playback.withProbeBudget(probesize: probesize, maxAnalyzeDuration: maxAnalyzeDuration).withSequentialOrigin(sequentialOrigin, declaredDuration: declaredDuration).withHeldSourceConnection(heldSourceConnection)) } dem.onNetworkPhaseChanged = networkPhaseSink try await host.load( @@ -2233,10 +2237,12 @@ extension AetherEngine { // silently revert to the main title. Preopen the disc demuxer with the title so the selection // survives the reload (#67). Non-disc URL sources keep customPreopened nil and reopen by URL. let headers = loadedOptions.httpHeaders + let authorization = loadedOptions.httpRequestAuthorization do { customPreopened = try await Task.detached(priority: .userInitiated) { let d = Demuxer() - try d.open(url: url, extraHeaders: headers, profile: reloadProfile, selectTitleID: titleToReopen) + try d.open(url: url, extraHeaders: headers, requestAuthorization: authorization, + profile: reloadProfile, selectTitleID: titleToReopen) return d }.value } catch { @@ -2286,6 +2292,7 @@ extension AetherEngine { try await loadSoftware( url: url, sourceHTTPHeaders: loadedOptions.httpHeaders, + sourceHTTPAuthorization: loadedOptions.httpRequestAuthorization, startPosition: LiveReloadPolicy.resumePosition( isLive: loadedOptions.isLive, currentTime: resumeAt), audioSourceStreamIndex: audioStreamIndex, @@ -2348,6 +2355,7 @@ extension AetherEngine { try await loadNative( url: url, sourceHTTPHeaders: loadedOptions.httpHeaders, + sourceHTTPAuthorization: loadedOptions.httpRequestAuthorization, // Live rejoins at the live edge (see loadSoftware above). startPosition: LiveReloadPolicy.resumePosition( isLive: loadedOptions.isLive, currentTime: resumeAt), diff --git a/Sources/AetherEngine/AetherEngine+Subtitles.swift b/Sources/AetherEngine/AetherEngine+Subtitles.swift index 896f79350..18669c748 100644 --- a/Sources/AetherEngine/AetherEngine+Subtitles.swift +++ b/Sources/AetherEngine/AetherEngine+Subtitles.swift @@ -657,6 +657,7 @@ extension AetherEngine { let isCustom = isCustomSource if isCustom, customReader == nil { return } let headers = loadedOptions.httpHeaders + let authorization = loadedOptions.httpRequestAuthorization let formatHint = customFormatHint let probesize = loadedOptions.probesize let maxAnalyzeDuration = loadedOptions.maxAnalyzeDuration @@ -708,7 +709,7 @@ extension AetherEngine { guard let self else { return } let outcome = await self.runSubtitleForwardPrefetchSession( url: url, reader: attemptReader, formatHint: formatHint, headers: headers, - startAt: resumeAt, callerProbesize: probesize, + authorization: authorization, startAt: resumeAt, callerProbesize: probesize, callerMaxAnalyzeDuration: maxAnalyzeDuration, selectTitleID: titleID, store: store, leadSeconds: lead, link: link) guard outcome.exit.isRestartable, !Task.isCancelled else { return } @@ -760,6 +761,7 @@ extension AetherEngine { /// packets to the SubtitlePacketStore instead of decoded cues to native stores. nonisolated private func runSubtitleForwardPrefetchSession( url: URL, reader: IOReader?, formatHint: String?, headers: [String: String], + authorization: HTTPRequestAuthorization?, startAt: Double, callerProbesize: Int64?, callerMaxAnalyzeDuration: Int64?, selectTitleID: Int?, store: SubtitlePacketStore, leadSeconds: Double, link: SideReaderLinkArbiter? @@ -813,8 +815,8 @@ extension AetherEngine { try demuxer.open(reader: reader, formatHint: formatHint, profile: openProfile, selectTitleID: selectTitleID, discCacheKey: url.absoluteString) } else { - try demuxer.open(url: url, extraHeaders: headers, profile: openProfile, - selectTitleID: selectTitleID) + try demuxer.open(url: url, extraHeaders: headers, requestAuthorization: authorization, + profile: openProfile, selectTitleID: selectTitleID) } } catch { EngineLog.emit("[AetherEngine] #151 forward prefetch open failed: \(error)", category: .engine) @@ -1678,6 +1680,7 @@ extension AetherEngine { customClone = clone } let headers = loadedOptions.httpHeaders + let authorization = loadedOptions.httpRequestAuthorization let formatHint = customFormatHint let w = sourceVideoWidth > 0 ? sourceVideoWidth : 1920 let h = sourceVideoHeight > 0 ? sourceVideoHeight : 1080 @@ -1691,7 +1694,7 @@ extension AetherEngine { nativeSubtitleReadersTask = Task.detached(priority: .utility) { [weak self] in await self?.runNativeSubtitleReaders( url: url, reader: reader, formatHint: formatHint, headers: headers, - pairs: pairs, startAt: startAt, videoWidth: w, videoHeight: h, + authorization: authorization, pairs: pairs, startAt: startAt, videoWidth: w, videoHeight: h, callerProbesize: probesize, callerMaxAnalyzeDuration: maxAnalyzeDuration, selectTitleID: titleID, readToEOF: readToEOF, link: link ) @@ -1749,6 +1752,7 @@ extension AetherEngine { nonisolated private func runNativeSubtitleReaders( url: URL, reader: IOReader?, formatHint: String?, headers: [String: String], + authorization: HTTPRequestAuthorization? = nil, pairs: [(streamIndex: Int32, store: NativeSubtitleCueStore)], startAt: Double, videoWidth: Int32, videoHeight: Int32, callerProbesize: Int64? = nil, callerMaxAnalyzeDuration: Int64? = nil, @@ -1778,7 +1782,8 @@ extension AetherEngine { if let reader = reader { try demuxer.open(reader: reader, formatHint: formatHint, profile: openProfile, selectTitleID: selectTitleID, discCacheKey: url.absoluteString) } else { - try demuxer.open(url: url, extraHeaders: headers, profile: openProfile, selectTitleID: selectTitleID) + try demuxer.open(url: url, extraHeaders: headers, requestAuthorization: authorization, + profile: openProfile, selectTitleID: selectTitleID) } } catch { EngineLog.emit("[AetherEngine] native subtitle readers open failed: \(error)", category: .engine) diff --git a/Sources/AetherEngine/AetherEngine.swift b/Sources/AetherEngine/AetherEngine.swift index 03aecd5f4..9c6950fcd 100644 --- a/Sources/AetherEngine/AetherEngine.swift +++ b/Sources/AetherEngine/AetherEngine.swift @@ -3908,7 +3908,8 @@ public final class AetherEngine: ObservableObject { case .url(let u): // isLive configures the AVIOReader for endless-feed mode; must be set at open time because // the probe demuxer is reused as the session demuxer (avformat_open_input runs only once). - try probe.open(url: u, extraHeaders: options.httpHeaders, profile: probeProfile, isLive: options.isLive, selectTitleID: discTitleID) + try probe.open(url: u, extraHeaders: options.httpHeaders, + requestAuthorization: options.httpRequestAuthorization, profile: probeProfile, isLive: options.isLive, selectTitleID: discTitleID) case .custom(let reader, let formatHint): // isLive suppresses SEEK_END duration estimate on forward-only live readers; same open-time requirement. try probe.open(reader: reader, formatHint: formatHint, profile: probeProfile, isLive: options.isLive, selectTitleID: discTitleID) @@ -4210,6 +4211,7 @@ public final class AetherEngine: ObservableObject { try await loadAudio( url: url, sourceHTTPHeaders: options.httpHeaders, + sourceHTTPAuthorization: options.httpRequestAuthorization, startPosition: startPosition, audioSourceStreamIndex: resolvedInitialAudio >= 0 ? resolvedInitialAudio : nil, preopenedDemuxer: probeOpened ? probe : nil, @@ -4626,6 +4628,7 @@ public final class AetherEngine: ObservableObject { try await loadSoftware( url: url, sourceHTTPHeaders: options.httpHeaders, + sourceHTTPAuthorization: options.httpRequestAuthorization, startPosition: startPosition, audioSourceStreamIndex: selectedAudio, isLive: options.isLive, @@ -4675,6 +4678,7 @@ public final class AetherEngine: ObservableObject { try await loadNative( url: url, sourceHTTPHeaders: options.httpHeaders, + sourceHTTPAuthorization: options.httpRequestAuthorization, startPosition: startPosition, audioSourceStreamIndex: selectedAudio, keepDvh1TagWithoutDV: options.keepDvh1TagWithoutDV, diff --git a/Sources/AetherEngine/Demuxer/AVIOReader.swift b/Sources/AetherEngine/Demuxer/AVIOReader.swift index 522c46eaa..2558663b2 100644 --- a/Sources/AetherEngine/Demuxer/AVIOReader.swift +++ b/Sources/AetherEngine/Demuxer/AVIOReader.swift @@ -42,6 +42,11 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { private let url: URL private let extraHeaders: [String: String] + /// `LoadOptions.httpRequestAuthorization`: when set, its answer replaces `extraHeaders` on every + /// request this reader builds. Nil keeps the static headers. + private let authorizer: SourceRequestAuthorizer? + /// How long one request waits for the provider before it fails, as the HLS relay allows. + static let authorizationTimeoutDefault: TimeInterval = 10 /// #450: connections the BOUNDED pool may hold to one host. A throttle, and it is allowed to be /// one: every request on that pool ends (a 4 MB detour block, a size probe, the tail prefetch), /// so a request that waits here waits for one that is finishing. @@ -730,6 +735,18 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { private var connStatus = 0 // Retry-After seconds from a rate-limit status, honoured before reconnect. private var connRetryAfter: TimeInterval = 0 + /// Refreshable authorization on the pump: the headers the current generation carried (the + /// `rejectedHeaders` of a 401), the fresh set a 401 staged for its one retry, and where that + /// retry stands. winCond-guarded; a generation that delivers resets the recovery. + private var connSentHeaders: [String: String] = [:] + private var stagedAuthorizedHeaders: [String: String]? + private var authorizationRecovery = AuthorizationRecovery.ready + + /// `refused` latches until a generation delivers again. Its status is the 401 that stood, or 0 + /// when the provider itself refused or did not answer. + private enum AuthorizationRecovery: Equatable { + case ready, retried, refused(status: Int) + } // Bumped on every (re)connect; stale delegate callbacks are ignored. private var connGeneration = 0 /// #377: what is on the link for this reader. Either shape is ONE request against the @@ -1022,12 +1039,15 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { private let probeDrainLock = NSLock() private var drainingProbeRequest = false - init(url: URL, extraHeaders: [String: String] = [:], label: String = "source", chunkSize: Int = 4 * 1024 * 1024, prefetchEnabled: Bool = true, isLive: Bool = false, chunkRequestTimeout: TimeInterval = 35, chunkMaxRetries: Int = 3, boundedInitialFetch: Int64? = nil, sequentialOnly: Bool = false, connStallTimeout: TimeInterval = AVIOReader.connStallTimeoutDefault, windowHighWater: Int? = nil, heldConnection: Bool = false, probeControl: ProbeControl? = nil, probeRequestSession: URLSession? = nil) { + init(url: URL, extraHeaders: [String: String] = [:], requestAuthorization: HTTPRequestAuthorization? = nil, authorizationTimeout: TimeInterval = AVIOReader.authorizationTimeoutDefault, label: String = "source", chunkSize: Int = 4 * 1024 * 1024, prefetchEnabled: Bool = true, isLive: Bool = false, chunkRequestTimeout: TimeInterval = 35, chunkMaxRetries: Int = 3, boundedInitialFetch: Int64? = nil, sequentialOnly: Bool = false, connStallTimeout: TimeInterval = AVIOReader.connStallTimeoutDefault, windowHighWater: Int? = nil, heldConnection: Bool = false, probeControl: ProbeControl? = nil, probeRequestSession: URLSession? = nil) { self.probeControl = probeControl self.probeRequestSession = probeControl == nil ? nil : probeRequestSession self.url = url self.label = label self.extraHeaders = extraHeaders + self.authorizer = requestAuthorization.map { + SourceRequestAuthorizer($0, sourceURL: url, timeout: authorizationTimeout) + } self.chunkSize = chunkSize self.prefetchEnabled = prefetchEnabled self.isLive = isLive @@ -1070,15 +1090,26 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// arrived `auth=none`, and the post-seek request to the same pinned host 13 s later carried /// both `Authorization` and `X-Emby-Token`. One policy, applied where the request is built, so /// a pin cannot outflank it. - private func headers(for target: URL?) -> [String: String] { - RedirectHeaderPolicy.headersToReplay( - extraHeaders: extraHeaders, originalURL: url, redirectURL: target ?? url) + /// + /// With a provider, `base` is its answer for the source URL, asked once per request: the same + /// policy then decides where that answer may go, exactly as it does for static headers. A throw + /// means the provider refused, timed out or the reader closed, and the request must not be sent. + private func headers(for target: URL?, base: [String: String]? = nil) throws -> [String: String] { + let source = try base ?? authorizer?.headers() ?? extraHeaders + return RedirectHeaderPolicy.headersToReplay( + extraHeaders: source, originalURL: url, redirectURL: target ?? url) } - private func applyExtraHeaders(_ request: inout URLRequest) { - for (name, value) in headers(for: request.url) { + /// Sets the request's headers and returns them, so the delegate replays this request's own + /// headers on a redirect instead of asking the provider a second time. + @discardableResult + private func applyExtraHeaders(_ request: inout URLRequest, + base: [String: String]? = nil) throws -> [String: String] { + let headers = try headers(for: request.url, base: base) + for (name, value) in headers { request.setValue(value, forHTTPHeaderField: name) } + return headers } func open() throws { @@ -1118,8 +1149,12 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { try failIfStreamingRefused(fallbackStatus: 0) } else if prefetchEnabled { // #281: the parse seeks that follow this open are what the retained head exists for. + // One provider answer serves the open's first two requests, so a provider that does not + // answer costs the open one bounded wait rather than two. + let openHeaders = try authorizeOpen() winCond.lock() openPhaseActive = true + stagedAuthorizedHeaders = openHeaders winCond.unlock() // #551: bytes a host warmed for this source before anything asked to play it. When // they are here, this open owes the origin nothing for the span they cover. @@ -1129,7 +1164,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // therefore needs nothing this open has learned yet. A warm that already carries the // trailing object has no use for it. if warm?.tail == nil { - startTailPrefetch() + startTailPrefetch(headers: openHeaders) } // Playback path. The persistent connection's `Range: bytes=0-` request is itself // the size probe: its 206 Content-Range is folded into fileSize by @@ -1152,7 +1187,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { openPrefix = Array(warm.head.data.prefix(16)) } else { startPersistentConnection(at: 0, boundedTo: boundedInitialFetch) - gotData = awaitFirstPersistentData() + gotData = try awaitFirstAuthorizedData() openPrefix = firstWindowPrefix() } // AE#140: an HLS playlist URL misrouted onto the raw-byte live path. A live origin serves the @@ -1296,6 +1331,55 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return gotData } + /// The provider's answer for the open, nil without a provider. A refusal or timeout fails the + /// open before anything is sent. + private func authorizeOpen() throws -> [String: String]? { + guard let authorizer else { return nil } + do { + return try authorizer.headers() + } catch { + latchProviderRefusal() + throw openFailureForAuthorization(status: 0) + } + } + + /// `awaitFirstPersistentData` plus the open's share of the 401 recovery: one retry at byte 0 + /// with refreshed headers, then a typed failure once the provider's answer stands. + private func awaitFirstAuthorizedData() throws -> Bool { + while !awaitFirstPersistentData() { + winCond.lock() + let status = connEnded ? connStatus : 0 + winCond.unlock() + switch recoverAuthorization(status: status) { + case .notApplicable: + return false + case .retry: + startPersistentConnection(at: 0, boundedTo: boundedInitialFetch) + case .refused(let refusedStatus): + throw openFailureForAuthorization(status: refusedStatus) + } + } + return true + } + + /// Closes the reader and returns the open's typed failure: a 401 the provider could not answer + /// as that status, a provider that refused or did not answer as `authorizationUnavailable`. + private func openFailureForAuthorization(status: Int) -> AVIOReaderError { + EngineLog.emit( + "[AVIOReader] \(label) source authorization refused (status=\(status)); failing the open typed", + category: .demux) + markClosed() + close() + return status == 0 ? .authorizationUnavailable : .httpStatus(status) + } + + /// The provider refused or did not answer. Latches until a generation delivers again. + private func latchProviderRefusal() { + winCond.lock() + authorizationRecovery = .refused(status: 0) + winCond.unlock() + } + /// The streaming GET was answered with a status instead of a body, or the ranged open was /// already refused with one and the unranged GET then delivered nothing either. Either way the /// demuxer would be handed an empty stream (or an error page) and report it as invalid data; @@ -1340,6 +1424,12 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { close() throw AVIOReaderError.transportSecurityFailed(code: tlsCode) } + winCond.lock() + let authorization = authorizationRecovery + winCond.unlock() + if case .refused(let refusedStatus) = authorization, status == 0 { + throw openFailureForAuthorization(status: refusedStatus) + } guard status != 0 else { return } EngineLog.emit( "[AVIOReader] \(label) source refused: HTTP \(status); failing the open typed", @@ -1463,6 +1553,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// Must be called BEFORE acquiring the demuxer's access lock. func markClosed() { isClosed = true + authorizer?.cancel() // Wake any semaphore waits so the read callbacks can exit prefetchReady.signal() streamDataReady.signal() @@ -1507,6 +1598,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { guard !isFullyClosed else { return } isFullyClosed = true isClosed = true + authorizer?.cancel() if let ctx = context { // avio_context_free does NOT free ctx->buffer (verified, aviobuf.c). // Free ctx.pointee.buffer, not original av_malloc ptr: FFmpeg can @@ -2103,10 +2195,21 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { winCond.broadcast() winCond.unlock() if let refillFrom { - if !refillFaulted || chargeFaultedRunwayRefill(at: refillFrom, ahead: undrained, - status: refillStatus, - retryAfter: refillRetryAfter) { + let authorization = refillFaulted ? recoverAuthorization(status: refillStatus) : .notApplicable + switch authorization { + case .retry: timedReconnect(seek: false, at: refillFrom) + case .refused: + // Serve what is resident; the empty-window path fails the read. + winCond.lock() + nextFaultedRefillAt = .distantFuture + winCond.unlock() + case .notApplicable: + if !refillFaulted || chargeFaultedRunwayRefill(at: refillFrom, ahead: undrained, + status: refillStatus, + retryAfter: refillRetryAfter) { + timedReconnect(seek: false, at: refillFrom) + } } } continue @@ -2199,6 +2302,17 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { timedReconnect(seek: false, at: frontier) continue } + switch recoverAuthorization(status: status) { + case .retry: + timedReconnect(seek: false, at: frontier) + continue + case .refused(let refusedStatus): + EngineLog.emit("[AVIOReader] \(label) authorization refused at offset \(frontier) status=\(refusedStatus); failing the read", category: .demux) + emitNetworkPhase(.exhausted) // the provider has answered; reconnecting cannot help + return totalRead > 0 ? Int32(totalRead) : -1 + case .notApplicable: + break + } // A 429/503/509 is rate limiting, not a dead source: drive give-up + backoff off the // rate-limit streak, which (unlike unproductiveReconnects) survives the seekReconnect // that parse seeks fire, so a throttled origin fails cleanly instead of looping (#71). @@ -2407,6 +2521,55 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { } } + // MARK: - Refreshable authorization + + /// The provider refused, did not answer, or the reader closed before this generation's request + /// was built, so nothing went on the link. The generation ends at once and the refusal latches: + /// the read fails instead of running the reconnect ladder against a provider that has answered. + private func failGenerationForAuthorization(_ generation: Int, error: Error) { + winCond.lock() + if generation == connGeneration { + connEnded = true + connStatus = 0 + authorizationRecovery = .refused(status: 0) + } + winCond.broadcast() + winCond.unlock() + guard !isClosed else { return } + EngineLog.emit( + "[AVIOReader] \(label) gen=\(generation) request authorization failed: \(error.localizedDescription)", + category: .demux) + } + + private enum AuthorizationStep { case notApplicable, retry, refused(status: Int) } + + /// A 401 against a reader with a provider earns one retry at the same offset, and only when the + /// provider answers the rejected headers with a different Authorization value. Unchanged + /// credentials, a provider failure or a second 401 latch the refusal. Every other status, and a + /// reader without a provider, stays with the reconnect ladder. Demux-thread-only; the provider + /// wait is bounded by the authorizer's timeout. + private func recoverAuthorization(status: Int) -> AuthorizationStep { + guard let authorizer else { return .notApplicable } + winCond.lock() + let state = authorizationRecovery + let rejected = connSentHeaders + winCond.unlock() + if case .refused(let latched) = state { return .refused(status: latched) } + guard status == 401 else { return .notApplicable } + let fresh = state == .ready ? authorizer.refreshed(rejecting: rejected) : nil + winCond.lock() + if let fresh { + stagedAuthorizedHeaders = fresh + authorizationRecovery = .retried + } else { + authorizationRecovery = .refused(status: 401) + } + winCond.unlock() + guard fresh != nil else { return .refused(status: 401) } + EngineLog.emit("[AVIOReader] \(label) 401 with refreshed authorization; retrying once", category: .demux) + return .retry + } + /// Increments the consecutive rate-limited streak; returns true once the bounded cap is hit. /// Demux-thread-only. Deliberately NOT reset by `seekReconnect` (parse seeks must not mask a /// throttled origin into an endless reconnect loop, #71); only real read progress clears it. @@ -2484,9 +2647,9 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // #93/#96: a starved backward-scrub detour fetch must abort fast (the rescue reconnect serves // instantly), so this path uses the tight interactive budget, not the full chunk timeout. request.timeoutInterval = budget - applyExtraHeaders(&request) do { - let (data, response) = try syncRequest(request, budget: budget) + let sentHeaders = try applyExtraHeaders(&request) + let (data, response) = try syncRequest(request, headers: sentHeaders, budget: budget) if let http = response as? HTTPURLResponse { let status = http.statusCode if Self.isRateLimitStatus(status) { @@ -2594,7 +2757,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// body, so it is ALWAYS still in flight at that moment on any origin whose first byte costs /// anything. It never once served the read it exists for; it only added a request. Only a /// loopback origin, which answers before the race can be lost, made it look like it worked. - private func startTailPrefetch() { + private func startTailPrefetch(headers base: [String: String]?) { guard !isLive, !isClosed else { return } let url = requestURL() // #377: a speculative second request is the first thing to drop on an origin that allows @@ -2619,7 +2782,8 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { var request = URLRequest(url: url) request.setValue("bytes=-\(Self.tailPrefetchBytes)", forHTTPHeaderField: "Range") request.timeoutInterval = Self.effectiveDetourBudget(chunkRequestTimeout: chunkRequestTimeout) - applyExtraHeaders(&request) + // Speculative: a provider that cannot answer costs this fetch, never the open. + guard let sentHeaders = try? applyExtraHeaders(&request, base: base) else { return } // A delegate rather than a completion handler, and the distinction is load-bearing: a // completion handler only fires once the body is in hand, so an origin that does not @@ -2641,7 +2805,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { let delegate = TailPrefetchDelegate( expectedLength: Self.tailPrefetchBytes, - extraHeaders: headers(for: request.url) + extraHeaders: sentHeaders ) // #281 retest: one line per open, and the line the field needs. The advertised way to check // this fix was "does a bytes=-65536 request show up", which the engine never printed, so a @@ -2901,6 +3065,11 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { if isClosed { return } + winCond.lock() + let staged = stagedAuthorizedHeaders + stagedAuthorizedHeaders = nil + winCond.unlock() + var request = URLRequest(url: requestURL()) // #93 residual: a bounded open connection asks for a finite range so an origin that dribbles // the open-ended `bytes=0-` stream serves it as a fast finite GET. The 206 Content-Range still @@ -2921,7 +3090,14 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { request.setValue("bytes=\(askAsJoin ? 0 : offset)-", forHTTPHeaderField: "Range") } request.timeoutInterval = 0 // long-lived; stalls handled by the reader - applyExtraHeaders(&request) + // Asked before the origin slot is taken, so a slow provider never holds one. + let sentHeaders: [String: String] + do { + sentHeaders = try applyExtraHeaders(&request, base: staged) + } catch { + failGenerationForAuthorization(generation, error: error) + return + } // #377: take the origin slot before the connection goes on the link. The pump is the one // path that must never be refused a slot for long: it is the main line, and everything @@ -2936,7 +3112,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { transfer = HeldSourceConnection( url: request.url ?? requestURLForBudget, offset: offset, - extraHeaders: headers(for: request.url ?? requestURLForBudget), + extraHeaders: sentHeaders, userAgent: nil, label: label, generation: generation, @@ -2947,7 +3123,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { let delegate = PersistentReadDelegate( reader: self, generation: generation, - extraHeaders: headers(for: request.url), + extraHeaders: sentHeaders, ticket: ticket, originURL: requestURLForBudget ) @@ -2965,6 +3141,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return } activeTransfer = transfer + connSentHeaders = sentHeaders winCond.unlock() transfer.startTransfer() @@ -3132,6 +3309,8 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // a give-up latch — an origin that recovered after the faulted ladder capped out may // fault again later and deserves a fresh ladder, not `.distantFuture` forever). nextFaultedRefillAt = .distantPast + // The credential this generation carried is accepted, so a later 401 earns a new retry. + authorizationRecovery = .ready } let count = data.count // #310: delivery that lands with the backpressure end ALREADY recorded, i.e. after our @@ -3397,7 +3576,17 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { private func streamDownloadSync() { var request = URLRequest(url: url) request.timeoutInterval = 0 // No timeout for live streams - applyExtraHeaders(&request) + let sentHeaders: [String: String] + do { + sentHeaders = try applyExtraHeaders(&request) + } catch { + latchProviderRefusal() + streamLock.lock() + streamEnded = true + streamLock.unlock() + streamDataReady.signal() + return + } // #377: this connection is open for the whole session, so it holds its slot for the whole // session, which is exactly what it costs the origin. Scoped to this function because the @@ -3420,7 +3609,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { let semaphore = DispatchSemaphore(value: 0) let delegate = StreamingDelegate( - extraHeaders: headers(for: request.url), + extraHeaders: sentHeaders, onResponse: { [weak self] response in // Advisory length for the sequential-origin EOF/EIO distinction; -1 (chunked / // unknown) leaves the clean-end path as the only EOF source. @@ -3764,7 +3953,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { var request = URLRequest(url: url) request.setValue(range, forHTTPHeaderField: "Range") request.timeoutInterval = 20 - applyExtraHeaders(&request) + guard let sentHeaders = try? applyExtraHeaders(&request) else { return nil } // #377: `probeSession` runs on `URLSessionConfiguration.default`, so its own cap is 6 and // it composes with nothing. The staggered fan fires two fallbacks at once by design, which @@ -3781,7 +3970,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { } defer { OriginRequestBudget.shared.release(ticket) } - let delegate = ProbeDelegate(extraHeaders: headers(for: request.url)) + let delegate = ProbeDelegate(extraHeaders: sentHeaders) let task = (probeRequestSession ?? Self.probeSession).dataTask(with: request) task.delegate = delegate @@ -3817,12 +4006,12 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { var request = URLRequest(url: url) request.httpMethod = "HEAD" request.timeoutInterval = 5 - applyExtraHeaders(&request) do { + let sentHeaders = try applyExtraHeaders(&request) // Honour the still budget here too so the open-time HEAD fallback can't // ride the default 35s on a stalled origin during a cold/reopen scrub (#27). - let (_, response) = try syncRequest(request, budget: chunkRequestTimeout) + let (_, response) = try syncRequest(request, headers: sentHeaders, budget: chunkRequestTimeout) guard let http = response as? HTTPURLResponse, (200...299).contains(http.statusCode) else { let status = (response as? HTTPURLResponse)?.statusCode ?? -1 @@ -3855,21 +4044,29 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return nil } - private func fetchChunkAttempt(from offset: Int64, size: Int, forceSource: Bool) -> Data? { + /// `refreshedHeaders` is the provider's answer to a 401 on this same range: set, this is the one + /// retry it permits, so a second 401 is final. + private func fetchChunkAttempt(from offset: Int64, size: Int, forceSource: Bool, + refreshedHeaders: [String: String]? = nil) -> Data? { let usingCachedURL = !forceSource && cachedResolvedURL() != nil let target = forceSource ? url : requestURL() let rangeEnd = offset + Int64(size) - 1 var request = URLRequest(url: target) request.setValue("bytes=\(offset)-\(rangeEnd)", forHTTPHeaderField: "Range") request.timeoutInterval = min(15, chunkRequestTimeout) - applyExtraHeaders(&request) + guard let sentHeaders = try? applyExtraHeaders(&request, base: refreshedHeaders) else { return nil } var lastError: Error? for attempt in 0.. (Data, URLResponse) { + /// `headers` are the ones `applyExtraHeaders` set on `request`, replayed on a redirect. + private func syncRequest(_ request: URLRequest, headers: [String: String], + budget: TimeInterval = 35) throws -> (Data, URLResponse) { // #377: every short fetch the reader makes (detour blocks, size probes, HEAD) funnels // through here, so this is the one place that has to take an origin slot for all of them. // Scoped to the call: unlike the pump's, this request's life IS this function's. @@ -4059,7 +4258,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { for: slotURL, label: "\(label) fetch", timeout: Self.shortFetchSlotWaitSeconds) defer { OriginRequestBudget.shared.release(ticket) } - let delegate = ChunkFetchDelegate(extraHeaders: headers(for: request.url), + let delegate = ChunkFetchDelegate(extraHeaders: headers, bodyLimit: Self.expectedBodyBytes(for: request)) let task = (probeRequestSession ?? Self.chunkSession).dataTask(with: request) task.delegate = delegate @@ -4941,6 +5140,9 @@ enum AVIOReaderError: Error, Equatable, CustomStringConvertible, LocalizedError /// the same reason `httpStatus` is: without it the open surfaces FFmpeg's invalid data and a /// self-signed origin reads as a corrupt file. case transportSecurityFailed(code: Int) + /// `LoadOptions.httpRequestAuthorization` refused the source or did not answer in time, so no + /// request was sent. + case authorizationUnavailable var description: String { switch self { @@ -4952,6 +5154,8 @@ enum AVIOReaderError: Error, Equatable, CustomStringConvertible, LocalizedError case .httpStatus(let status): return "Origin answered HTTP \(status) for the source" case .transportSecurityFailed(let code): return TransportSecurityFailure.sentence(for: code) + case .authorizationUnavailable: + return "The request authorization provider supplied no headers for the source" } } diff --git a/Sources/AetherEngine/Demuxer/Demuxer.swift b/Sources/AetherEngine/Demuxer/Demuxer.swift index 7465a04b1..78e87629d 100644 --- a/Sources/AetherEngine/Demuxer/Demuxer.swift +++ b/Sources/AetherEngine/Demuxer/Demuxer.swift @@ -425,14 +425,17 @@ public final class Demuxer: @unchecked Sendable { /// Open a media URL and probe its streams. /// - Parameters: /// - extraHeaders: Attached to every HTTP request (ignored for file:// URLs). + /// - requestAuthorization: `LoadOptions.httpRequestAuthorization`. Replaces `extraHeaders` on + /// the byte-range reader's requests; a remote disc image keeps `extraHeaders`. /// - isLive: Suppresses EOF synthesis and surfaces terminal error on reconnect cap. - func open(url: URL, extraHeaders: [String: String] = [:], profile: DemuxerOpenProfile = .playback, isLive: Bool = false, selectTitleID: Int? = nil) throws { + func open(url: URL, extraHeaders: [String: String] = [:], requestAuthorization: HTTPRequestAuthorization? = nil, profile: DemuxerOpenProfile = .playback, isLive: Bool = false, selectTitleID: Int? = nil) throws { self.openProfile = profile self.auditSource = isLive ? nil : (url, extraHeaders) let isHTTP = url.scheme == "http" || url.scheme == "https" if isHTTP { - try openHTTP(url: url, extraHeaders: extraHeaders, isLive: isLive, selectTitleID: selectTitleID) + try openHTTP(url: url, extraHeaders: extraHeaders, requestAuthorization: requestAuthorization, + isLive: isLive, selectTitleID: selectTitleID) } else { // Route a local DVD ISO through the disc adapter (FileIOReader keeps it // out of RAM). Falls back to the normal local open when not a disc. @@ -502,7 +505,7 @@ public final class Demuxer: @unchecked Sendable { ["iso", "img", "udf"].contains(url.pathExtension.lowercased()) } - private func openHTTP(url: URL, extraHeaders: [String: String], isLive: Bool = false, selectTitleID: Int? = nil) throws { + private func openHTTP(url: URL, extraHeaders: [String: String], requestAuthorization: HTTPRequestAuthorization?, isLive: Bool = false, selectTitleID: Int? = nil) throws { // A remote disc image goes through the same disc adapter as a local ISO (a raw .iso handed // straight to libavformat fails to probe; it is a filesystem, not a media container, #64). // Gated on the disc-image extension so normal media URLs skip the range-probe entirely; if @@ -522,6 +525,7 @@ public final class Demuxer: @unchecked Sendable { let reader = AVIOReader( url: url, extraHeaders: extraHeaders, + requestAuthorization: requestAuthorization, label: openProfile.readerLabel, chunkSize: openProfile.avioChunkSize, prefetchEnabled: openProfile.avioPrefetch, diff --git a/Sources/AetherEngine/Network/HLSOriginRelay.swift b/Sources/AetherEngine/Network/HLSOriginRelay.swift index 534adabdf..33f3d5b16 100644 --- a/Sources/AetherEngine/Network/HLSOriginRelay.swift +++ b/Sources/AetherEngine/Network/HLSOriginRelay.swift @@ -413,7 +413,7 @@ final class HLSOriginRelay: @unchecked Sendable { for (key, value) in resolved { // Header ownership is independent of case. The application cannot change byte // selection, routing, or framing by returning a header dictionary. - if authorization != nil && Self.transportHeaders.contains(key.lowercased()) { continue } + if authorization != nil && HTTPRequestAuthorization.transportHeaders.contains(key.lowercased()) { continue } request.setValue(value, forHTTPHeaderField: key) } if let range { request.setValue(range, forHTTPHeaderField: "Range") } @@ -449,7 +449,7 @@ final class HLSOriginRelay: @unchecked Sendable { let respondingURL = response.url ?? url let fresh = (try? authorize(respondingURL, rejectedHeaders: sentHeaders, fallback: [:])) .map { downgradeSafe($0, to: respondingURL) } - guard let fresh, Self.authorizationValue(fresh) != Self.authorizationValue(sentHeaders) else { + guard let fresh, HTTPRequestAuthorization.authorizationValue(fresh) != HTTPRequestAuthorization.authorizationValue(sentHeaders) else { reportRequestFailure() return .held(Fetched(url: respondingURL, status: 401, body: Data(), contentType: nil, contentRange: nil)) } @@ -459,16 +459,8 @@ final class HLSOriginRelay: @unchecked Sendable { } } - private static let transportHeaders: Set = [ - "range", "host", "content-length", "transfer-encoding", "connection", "trailer", "te", "upgrade" - ] - private static func isTLS(_ url: URL) -> Bool { url.scheme?.lowercased() == "https" } - private static func authorizationValue(_ headers: [String: String]) -> String? { - headers.first { $0.key.caseInsensitiveCompare("Authorization") == .orderedSame }?.value - } - private func remainingBudget(upTo maximum: TimeInterval) -> TimeInterval { guard let deadline else { return maximum } return max(0, min(maximum, deadline.timeIntervalSinceNow)) diff --git a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift index 22d5fb1e6..1252d836b 100644 --- a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift +++ b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift @@ -2,10 +2,15 @@ import Foundation /// Supplies the complete application headers for engine-owned HTTP requests. /// -/// Set `LoadOptions.httpRequestAuthorization` to keep a native HLS item playing as credentials -/// change. Set `ExternalSubtitleTrack.httpRequestAuthorization` separately for sidecar requests; -/// `data(from:maximumBytes:)` fetches bounded auxiliary resources with the same transport policy. -/// Direct media AVIO and live ingest still use their existing static headers. +/// Set `LoadOptions.httpRequestAuthorization` to keep a native HLS item or a direct-play source +/// playing as credentials change. Set `ExternalSubtitleTrack.httpRequestAuthorization` separately +/// for sidecar requests; `data(from:maximumBytes:)` fetches bounded auxiliary resources with the same +/// transport policy. Live ingest still uses its static headers. +/// +/// Direct media (the engine's own byte-range reader) asks for the source URL the host loaded before +/// every request it builds: each range, reconnect, probe and seek. The answer then follows the +/// static-header redirect policy, so credentials reach only the source's origin (and an http-to-https +/// upgrade of it), never a cross-origin redirect target or a target pinned from one. /// The resolver must independently validate every URL, including redirects and playlist-discovered /// origins. Discovery grants no credential authority. Credentials must never be placed in URLs. /// Include scheme, host and effective port in that scope. Redirects are authorized afresh. Throw to @@ -38,6 +43,15 @@ public final class HTTPRequestAuthorization: Sendable, Equatable { static let resourceTransferTimeout: TimeInterval = 20 + /// Headers the engine owns. A provider's answer cannot change byte selection, routing or framing. + static let transportHeaders: Set = [ + "range", "host", "content-length", "transfer-encoding", "connection", "trailer", "te", "upgrade" + ] + + static func authorizationValue(_ headers: [String: String]) -> String? { + headers.first { $0.key.caseInsensitiveCompare("Authorization") == .orderedSame }?.value + } + public static func == (lhs: HTTPRequestAuthorization, rhs: HTTPRequestAuthorization) -> Bool { lhs === rhs } @@ -88,3 +102,61 @@ final class HTTPAuthorizationWait: @unchecked Sendable { running?.cancel() } } + +/// The provider as the byte-range reader sees it: synchronous, bounded, and always asked about the +/// source URL. The reader calls it from the demux thread and its probe threads, so every wait ends at +/// `timeout` or at `cancel()`, whichever comes first, even when the host never answers. +final class SourceRequestAuthorizer: @unchecked Sendable { + private let authorization: HTTPRequestAuthorization + private let sourceURL: URL + private let timeout: TimeInterval + private let lock = NSLock() + private var waits: [ObjectIdentifier: HTTPAuthorizationWait] = [:] + private var cancelled = false + + init(_ authorization: HTTPRequestAuthorization, sourceURL: URL, timeout: TimeInterval) { + self.authorization = authorization + self.sourceURL = sourceURL + self.timeout = timeout + } + + /// The complete headers for a new request, without the ones the engine owns. Throws when the + /// provider refuses, does not answer within `timeout`, or the reader has closed. + func headers(rejecting rejected: [String: String]? = nil) throws -> [String: String] { + let wait = HTTPAuthorizationWait() + let id = ObjectIdentifier(wait) + lock.lock() + guard !cancelled else { + lock.unlock() + throw CancellationError() + } + waits[id] = wait + lock.unlock() + defer { + lock.lock() + waits[id] = nil + lock.unlock() + } + let answer = try wait.resolve(authorization, url: sourceURL, rejectedHeaders: rejected, + timeout: timeout) + return answer.filter { !HTTPRequestAuthorization.transportHeaders.contains($0.key.lowercased()) } + } + + /// The answer to one HTTP 401 against `rejected`, the headers that request actually carried. Nil + /// unless the provider returns a different Authorization value, which is what permits one retry. + func refreshed(rejecting rejected: [String: String]) -> [String: String]? { + guard let fresh = try? headers(rejecting: rejected), + HTTPRequestAuthorization.authorizationValue(fresh) + != HTTPRequestAuthorization.authorizationValue(rejected) else { return nil } + return fresh + } + + /// Ends every pending wait and refuses later ones. Called when the reader closes. + func cancel() { + lock.lock() + cancelled = true + let pending = Array(waits.values) + lock.unlock() + pending.forEach { $0.cancel() } + } +} diff --git a/Sources/AetherEngine/PlayerState.swift b/Sources/AetherEngine/PlayerState.swift index c11d05a55..6f51a08d6 100644 --- a/Sources/AetherEngine/PlayerState.swift +++ b/Sources/AetherEngine/PlayerState.swift @@ -408,9 +408,11 @@ public struct LoadOptions: Sendable, Equatable { public var suppressDisplayCriteria: Bool /// Extra HTTP headers for HEAD probe, Range chunks, side-demuxer fetches. On the loopback paths they are NOT forwarded to AVPlayer (it hits the local server); on `nativeRemoteHLS` they ride into the AVURLAsset so header-enforcing origins (IPTV Referer / User-Agent / Authorization) work (#119). Forwarded to `selectSidecarSubtitle` by default; pass explicit headers to override (#32). Default empty. public var httpHeaders: [String: String] - /// Refreshable complete headers for native HLS relay requests and playlist preflight. Forces - /// engine-owned transport from the initial load. Does not apply to AVIO, live ingest, or external - /// subtitle downloads. On supported requests this replaces `httpHeaders`; default nil. + /// Refreshable complete headers for native HLS relay requests and playlist preflight, and for + /// every range request of a direct-play source (its reopens and embedded-subtitle side readers + /// included). Forces engine-owned transport for native HLS from the initial load. Does not apply + /// to live ingest, remote disc images, or external subtitle downloads. On supported requests + /// this replaces `httpHeaders`; default nil. public var httpRequestAuthorization: HTTPRequestAuthorization? /// Diagnostic lever: force dvh1 codec tags + master playlist regardless of display capability. OFF by default: non-DV displays route DV through the media playlist (no master) so AVPlayer auto-tonemaps the HEVC base layer (only path that avoids AVFoundationErrorDomain -11868 on tvOS 26). AetherEngine#4. diff --git a/Sources/AetherEngine/Video/HLSVideoEngine+LiveReopen.swift b/Sources/AetherEngine/Video/HLSVideoEngine+LiveReopen.swift index 617085ea9..c2016c671 100644 --- a/Sources/AetherEngine/Video/HLSVideoEngine+LiveReopen.swift +++ b/Sources/AetherEngine/Video/HLSVideoEngine+LiveReopen.swift @@ -909,7 +909,8 @@ extension HLSVideoEngine { do { switch transport { case .url: - try dem.open(url: sourceURL, extraHeaders: sourceHTTPHeaders, profile: openProfile, isLive: true) + try dem.open(url: sourceURL, extraHeaders: sourceHTTPHeaders, + requestAuthorization: sourceHTTPAuthorization, profile: openProfile, isLive: true) case .customFactory: // #199: fresh engine-created ingest reader over the same channel; the dead // reader's construction inputs are immutable, so this rejoins at the live edge. diff --git a/Sources/AetherEngine/Video/HLSVideoEngine.swift b/Sources/AetherEngine/Video/HLSVideoEngine.swift index 0fdbca27c..58133f26f 100644 --- a/Sources/AetherEngine/Video/HLSVideoEngine.swift +++ b/Sources/AetherEngine/Video/HLSVideoEngine.swift @@ -42,6 +42,8 @@ public final class HLSVideoEngine: @unchecked Sendable { let sourceURL: URL let sourceHTTPHeaders: [String: String] + /// `LoadOptions.httpRequestAuthorization`, carried to every reopen of the source. + let sourceHTTPAuthorization: HTTPRequestAuthorization? private let dvModeAvailable: Bool /// From `LoadOptions.keepDvh1TagWithoutDV`; default OFF, set only for misreporting DV panels. @@ -849,6 +851,7 @@ public final class HLSVideoEngine: @unchecked Sendable { public init( url: URL, sourceHTTPHeaders: [String: String] = [:], + sourceHTTPAuthorization: HTTPRequestAuthorization? = nil, dvModeAvailable: Bool = true, displaySupportsHDR: Bool = true, keepDvh1TagWithoutDV: Bool = false, @@ -882,6 +885,7 @@ public final class HLSVideoEngine: @unchecked Sendable { ) { self.sourceURL = url self.sourceHTTPHeaders = sourceHTTPHeaders + self.sourceHTTPAuthorization = sourceHTTPAuthorization self.sequentialOrigin = sequentialOrigin self.heldSourceConnection = heldSourceConnection self.declaredDurationSeconds = declaredDurationSeconds @@ -1091,7 +1095,8 @@ public final class HLSVideoEngine: @unchecked Sendable { } else { dem = Demuxer() do { - try dem.open(url: sourceURL, extraHeaders: sourceHTTPHeaders, profile: openProfile, isLive: isLiveSession) + try dem.open(url: sourceURL, extraHeaders: sourceHTTPHeaders, + requestAuthorization: sourceHTTPAuthorization, profile: openProfile, isLive: isLiveSession) } catch { throw Self.openFailure(from: error) } @@ -3876,6 +3881,7 @@ public final class HLSVideoEngine: @unchecked Sendable { // reopen would splice fabricated-position bytes into the new pump. try fresh.open( url: sourceURL, extraHeaders: sourceHTTPHeaders, + requestAuthorization: sourceHTTPAuthorization, profile: restartReopenProfile, isLive: false) dem.markClosed() // abort any wedged read now that the replacement is ready diff --git a/Tests/AetherEngineTests/Issue255BodyReserveTests.swift b/Tests/AetherEngineTests/Issue255BodyReserveTests.swift index bf140bd71..b0bc22c1a 100644 --- a/Tests/AetherEngineTests/Issue255BodyReserveTests.swift +++ b/Tests/AetherEngineTests/Issue255BodyReserveTests.swift @@ -27,6 +27,7 @@ final class ScriptedOriginServer: @unchecked Sendable { struct Recorded: Sendable { let method: String let range: String? + var authorization: String? = nil } struct Reply { @@ -162,7 +163,9 @@ final class ScriptedOriginServer: @unchecked Sendable { let method = lines.first?.split(separator: " ").first.map(String.init) ?? "GET" let range = lines.first(where: { $0.lowercased().hasPrefix("range:") }) .map { String($0.dropFirst("range:".count)).trimmingCharacters(in: .whitespaces) } - let recorded = Recorded(method: method, range: range) + let authorization = lines.first(where: { $0.lowercased().hasPrefix("authorization:") }) + .map { String($0.dropFirst("authorization:".count)).trimmingCharacters(in: .whitespaces) } + let recorded = Recorded(method: method, range: range, authorization: authorization) lock.lock() _requests.append(recorded) lock.unlock() diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift new file mode 100644 index 000000000..e91de622e --- /dev/null +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -0,0 +1,213 @@ +// `LoadOptions.httpRequestAuthorization` on direct play. The byte-range reader used to send the +// headers it was opened with for the whole session, so once a host rotated its access token every +// later range was refused and the host had to reload the player at the current position. The reader +// now asks the provider before each request it builds, gives a 401 one retry at the same offset when +// the provider answers with a different credential, and bounds every wait on the provider. +// +// Each test serves a 256 KiB source whose open is bounded to the first 64 KiB, so playback through it +// takes exactly two pump ranges: `bytes=0-65535`, then `bytes=65536-262143`. +import Foundation +import Testing +@testable import AetherEngine + +@Suite("Refreshable direct-play authorization", .timeLimit(.minutes(1))) +struct RefreshableDirectPlayAuthorizationTests { + private static let total = 256 * 1024 + private static let openBytes = 64 * 1024 + private static let secondRange = "bytes=65536-262143" + + @Test("every range request carries the provider's current answer, not the static headers") + func headersResolvedPerRange() async throws { + let provider = ProviderLog() + let server = try Self.origin { _ in true } + defer { server.stop() } + let reader = Self.reader(server, headers: ["Authorization": "Bearer static"], + authorization: HTTPRequestAuthorization { _, rejected in + provider.answer(rejected: rejected) { "Bearer \($0)" } + }) + defer { reader.markClosed(); reader.close() } + + let read = try await offThread(reader) { + try reader.open() + return Self.read(reader, count: Self.total) + } + + #expect(read.bytes == Self.total) + let first = try #require(server.requests.first { $0.range == "bytes=0-65535" }) + let second = try #require(server.requests.first { $0.range == Self.secondRange }) + #expect(first.authorization?.hasPrefix("Bearer ") == true) + #expect(second.authorization?.hasPrefix("Bearer ") == true) + #expect(first.authorization != second.authorization, "the second range reused the first answer") + #expect(server.requests.allSatisfy { $0.authorization != "Bearer static" }) + #expect(provider.rejections.isEmpty) + } + + @Test("a 401 is retried once at the same offset with the refreshed credential") + func unauthorizedRetriedOnceWithFreshCredential() async throws { + let provider = ProviderLog() + // The token rotates between the two ranges: the origin accepts only the fresh one past the open. + let server = try Self.origin { $0.range?.hasPrefix("bytes=65536") != true || $0.authorization == "Bearer fresh" } + defer { server.stop() } + let reader = Self.reader(server, authorization: HTTPRequestAuthorization { _, rejected in + provider.answer(rejected: rejected) { _ in rejected == nil && !provider.rotated ? "Bearer stale" : "Bearer fresh" } + }) + defer { reader.markClosed(); reader.close() } + + let read = try await offThread(reader) { + try reader.open() + return Self.read(reader, count: Self.total) + } + + #expect(read.bytes == Self.total, "playback stopped at \(read.bytes) with \(read.result)") + let atOffset = server.requests.filter { $0.range == Self.secondRange }.map(\.authorization) + #expect(atOffset == ["Bearer stale", "Bearer fresh"]) + #expect(provider.rejections == ["Bearer stale"]) + } + + @Test("an unchanged credential fails the read after one 401 instead of reconnecting") + func unchangedCredentialFailsTheRead() async throws { + let provider = ProviderLog() + let server = try Self.origin { $0.range?.hasPrefix("bytes=65536") != true } + defer { server.stop() } + let reader = Self.reader(server, authorization: HTTPRequestAuthorization { _, rejected in + provider.answer(rejected: rejected) { _ in "Bearer stale" } + }) + defer { reader.markClosed(); reader.close() } + + let read = try await offThread(reader) { + try reader.open() + return Self.read(reader, count: Self.total) + } + + #expect(read.bytes == Self.openBytes) + #expect(read.result < 0) + #expect(server.requests.filter { $0.range == Self.secondRange }.count == 1) + #expect(provider.rejections == ["Bearer stale"]) + } + + @Test("a provider that does not answer fails the open typed, and nothing is sent") + func providerTimeoutFailsTheOpen() async throws { + let server = try Self.origin { _ in true } + defer { server.stop() } + let reader = Self.reader(server, authorizationTimeout: 0.2, + authorization: HTTPRequestAuthorization { _, _ in + try await Task.sleep(for: .seconds(60)) + return [:] + }) + defer { reader.markClosed(); reader.close() } + + let started = Date() + let opened: Result = try await offThread(reader) { Result { try reader.open() } } + + #expect(throws: AVIOReaderError.authorizationUnavailable) { try opened.get() } + #expect(Date().timeIntervalSince(started) < 5) + #expect(server.requests.isEmpty) + } + + @Test("a provider that stops answering mid-stream fails the read within its bound") + func providerTimeoutFailsTheRead() async throws { + let provider = ProviderLog() + let server = try Self.origin { _ in true } + defer { server.stop() } + let reader = Self.reader(server, authorizationTimeout: 0.2, + authorization: HTTPRequestAuthorization { _, rejected in + if provider.rotated { try await Task.sleep(for: .seconds(60)) } + return provider.answer(rejected: rejected) { "Bearer \($0)" } + }) + defer { reader.markClosed(); reader.close() } + + try await offThread(reader) { try reader.open() } + provider.rotate() + let started = Date() + let read = try await offThread(reader) { Self.read(reader, count: Self.total) } + + #expect(read.bytes == Self.openBytes) + #expect(read.result < 0) + #expect(Date().timeIntervalSince(started) < 5) + #expect(!server.requests.contains { $0.range == Self.secondRange }) + } + + // MARK: - Support + + /// A ranged origin over `total` filler bytes. `accepts` decides per request; a refusal is a 401. + private static func origin(accepts: @escaping @Sendable (ScriptedOriginServer.Recorded) -> Bool) throws + -> ScriptedOriginServer { + try #require(ScriptedOriginServer { request in + guard accepts(request) else { return .init(status: 401, declaredLength: 0) } + let (start, end) = Self.bounds(request.range) + return .init(status: 206, declaredLength: Int64(end - start + 1), + contentRange: "bytes \(start)-\(end)/\(total)", bodyBytes: end - start + 1) + }) + } + + private static func bounds(_ range: String?) -> (Int, Int) { + let spec = range.map { String($0.dropFirst("bytes=".count)) } ?? "0-" + if spec.hasPrefix("-") { return (total - (Int(spec.dropFirst()) ?? 0), total - 1) } + let parts = spec.split(separator: "-", omittingEmptySubsequences: false) + let start = Int(parts[0]) ?? 0 + let end = parts.count > 1 ? Int(parts[1]).map { min($0, total - 1) } ?? total - 1 : total - 1 + return (start, end) + } + + private static func reader(_ server: ScriptedOriginServer, headers: [String: String] = [:], + authorizationTimeout: TimeInterval = AVIOReader.authorizationTimeoutDefault, + authorization: HTTPRequestAuthorization) -> AVIOReader { + AVIOReader(url: URL(string: "http://127.0.0.1:\(server.port)/movie.mkv")!, + extraHeaders: headers, requestAuthorization: authorization, + authorizationTimeout: authorizationTimeout, + boundedInitialFetch: Int64(openBytes)) + } + + /// Reads until `count` bytes or the first result that is not a delivery. + private static func read(_ reader: AVIOReader, count: Int) -> (bytes: Int, result: Int32) { + var buffer = [UInt8](repeating: 0, count: 64 * 1024) + var bytes = 0 + var result: Int32 = 0 + while bytes < count { + result = buffer.withUnsafeMutableBufferPointer { + reader.read(into: $0.baseAddress!, size: Int32(min($0.count, count - bytes))) + } + guard result > 0 else { break } + bytes += Int(result) + } + return (bytes, result) + } + + /// The reader's calls block, so they run on their own thread; cancelling closes the reader, which + /// ends whatever wait the thread is in. + private func offThread(_ reader: AVIOReader, + _ body: @escaping @Sendable () throws -> T) async throws -> T { + try await withTaskCancellationHandler { + try await withCheckedThrowingContinuation { continuation in + Thread.detachNewThread { continuation.resume(with: Result { try body() }) } + } + } onCancel: { + reader.markClosed() + } + } +} + +/// What the provider was asked. `answer` numbers every call from 1 and records the Authorization of +/// each set of rejected headers it was handed. +private final class ProviderLog: @unchecked Sendable { + private let lock = NSLock() + private var calls = 0 + private var _rejections: [String?] = [] + private var _rotated = false + + var rejections: [String?] { lock.withLock { _rejections } } + var rotated: Bool { lock.withLock { _rotated } } + func rotate() { lock.withLock { _rotated = true } } + + func answer(rejected: [String: String]?, token: (Int) -> String) -> [String: String] { + lock.lock() + calls += 1 + let call = calls + if let rejected { + _rejections.append(HTTPRequestAuthorization.authorizationValue(rejected)) + _rotated = true + } + lock.unlock() + return ["Authorization": token(call)] + } +} diff --git a/docs/api.md b/docs/api.md index 4d04f8d47..44a8a3982 100644 --- a/docs/api.md +++ b/docs/api.md @@ -167,7 +167,19 @@ substitute for scope checks: credentials in a custom header are not recognized, starts on HTTP sends whatever the resolver returns. `LoadOptions.httpRequestAuthorization` covers native HLS media and its master/variant playlist -preparation. Direct media AVIO, live ingest and audio taps retain static headers. +preparation, and direct play through the engine's byte-range reader: the playback source, its +reopens and reloads, and the side readers that pull embedded subtitles from it. Live ingest, audio +taps, remote disc images, one-shot probes and scrub thumbnails retain static headers. + +On direct play the reader asks the resolver for the source URL the host loaded before every +request it builds: each range, reconnect, seek, size probe and tail fetch. The answer replaces +`httpHeaders` and then follows the static-header redirect policy, so credentials reach only the +source's origin (or an http-to-https upgrade of it), never a cross-origin redirect target or a +target pinned from one. The resolver is not asked about those destinations. After a 401 the +resolver receives the headers that request carried; a changed `Authorization` value retries the +request once at the same byte offset. Unchanged credentials, a second 401, or a resolver that throws +or does not answer within its bound end the read instead of running the reconnect ladder, and fail +an open before anything is sent. A rotated token therefore needs no player reload. For external subtitles, set `ExternalSubtitleTrack.httpRequestAuthorization` on each registered track. This is independent of the media provider, so the host can restrict subtitle credentials @@ -967,7 +979,7 @@ All flags default to safe values; the table is the full set. Depth for the media | Option | Default | What it does | | --- | --- | --- | | `httpHeaders` | empty | Static headers on probes, range and segment fetches. On direct `nativeRemoteHLS` loads they ride into the `AVURLAsset`; when relayed they stay on upstream requests. Forwarded to sidecar subtitle fetches unless overridden. | -| `httpRequestAuthorization` | nil | An `HTTPRequestAuthorization` resolver for native HLS requests and playlist preparation. Forces the engine relay, replaces static application headers per request, and supports a bounded changed-bearer retry after 401. See [the credential contract](#rotating-native-hls-credentials-without-replacing-the-item). | +| `httpRequestAuthorization` | nil | An `HTTPRequestAuthorization` resolver for native HLS requests and playlist preparation, and for every range request of a direct-play source. Replaces static application headers per request and supports a bounded changed-bearer retry after 401; native HLS also forces the engine relay. See [the credential contract](#rotating-native-hls-credentials-without-replacing-the-item). | | `isLive` | false | Treat the source as live. Set it explicitly; duration-based auto-detection is too noisy. | | `dvrWindowSeconds` | nil | Timeshift window. nil means live-only and `seek` is a no-op. | | `liveJoinProfile` | `.standard` | A `LiveJoinProfile`. `.fastZap` collapses TARGETDURATION to the source GOP so an IPTV join costs seconds instead of a full holdback. | From ebc1129a1fb111f8d10ecdfcd3d2d5d6c5ef7b67 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:03:38 -0400 Subject: [PATCH 2/9] docs(api): note that natively decoded audio-only sources keep static headers The direct-play authorization notes implied every progressive source follows the resolver. An audio-only source whose codec AVPlayer decodes is probed through the authorized reader, then handed to AVPlayer with the static httpHeaders, so token rotation does not reach it. Name that exception in the credential contract and the changelog. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- docs/api.md | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ba3830434..9f4ff0ad2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,7 @@ the public-API contract. - `ExternalSubtitleTrack.httpRequestAuthorization` supplies refreshable headers for primary/secondary sidecars and native subtitle stores without changing registered track IDs or rendition mappings. Authorized container decoding retains AVIO streaming and range access. - `HTTPRequestAuthorization.data(from:maximumBytes:)` fetches raw auxiliary resources such as font bundles with a caller-supplied byte limit and a whole-transfer deadline, reusing the relay's redirect, authorization, retry, cancellation and TLS policy. -- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest and remote disc images keep static headers. +- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. - `LoadOptions.httpRequestAuthorization` accepts an async `HTTPRequestAuthorization` resolver for native HLS. The engine resolves headers before requests and redirects, and retries a rejected request once when the bearer changes, preserving the active player item across token rotation. diff --git a/docs/api.md b/docs/api.md index 44a8a3982..60f48c195 100644 --- a/docs/api.md +++ b/docs/api.md @@ -169,7 +169,9 @@ starts on HTTP sends whatever the resolver returns. `LoadOptions.httpRequestAuthorization` covers native HLS media and its master/variant playlist preparation, and direct play through the engine's byte-range reader: the playback source, its reopens and reloads, and the side readers that pull embedded subtitles from it. Live ingest, audio -taps, remote disc images, one-shot probes and scrub thumbnails retain static headers. +taps, remote disc images, one-shot probes and scrub thumbnails retain static headers. So does an +audio-only progressive source whose codec AVPlayer decodes (AAC, MP3, MP1/MP2, FLAC, ALAC, AC-3, +E-AC-3, PCM): the resolver authorizes its probe, then AVPlayer plays it with `httpHeaders`. On direct play the reader asks the resolver for the source URL the host loaded before every request it builds: each range, reconnect, seek, size probe and tail fetch. The answer replaces From bf515f59d4ccd2684d22989788e5e1071dbcfe87 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 13:55:34 -0400 Subject: [PATCH 3/9] fix(auth): run request resolvers off the cooperative pool Every demuxer open runs inside a Task.detached, so with a direct-play resolver the open parks a cooperative-pool thread on an NSCondition while it waits for the resolver. The resolver itself ran in a Task.detached on that same pool. Opens that occupied every pool thread left their resolvers nowhere to start, so each failed at its 10 s bound as authorizationUnavailable. Before this branch the opens also blocked pool threads, but only on URLSession I/O, which completes on its own delegate queue and needs nothing from the pool. HTTPAuthorizationWait now starts the resolver with a task executor preference for a process-wide ResolverExecutor backed by a serial dispatch queue. A serial queue draws its thread from the overcommit root, so it starts while the pool is parked, and default actors the resolver awaits run there too. visionOS 1 keeps the previous detached task. The HLS relay shares the wait and gets the same isolation. The new test parks twice the pool's width of callers on one resolver that awaits an actor. Before the fix 22 of 32 timed out, and the parked pool starved the sibling tests in the suite. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- .../Network/HTTPRequestAuthorization.swift | 37 ++++++++++++++++--- ...reshableDirectPlayAuthorizationTests.swift | 28 ++++++++++++++ docs/api.md | 4 +- 4 files changed, 63 insertions(+), 8 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9f4ff0ad2..c92b6e4a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,7 @@ the public-API contract. - `ExternalSubtitleTrack.httpRequestAuthorization` supplies refreshable headers for primary/secondary sidecars and native subtitle stores without changing registered track IDs or rendition mappings. Authorized container decoding retains AVIO streaming and range access. - `HTTPRequestAuthorization.data(from:maximumBytes:)` fetches raw auxiliary resources such as font bundles with a caller-supplied byte limit and a whole-transfer deadline, reusing the relay's redirect, authorization, retry, cancellation and TLS policy. -- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. +- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. Resolvers run on an engine-owned serial executor, off Swift's cooperative pool, so demuxer opens that occupy every pool thread while they wait cannot starve the resolver they wait for. - `LoadOptions.httpRequestAuthorization` accepts an async `HTTPRequestAuthorization` resolver for native HLS. The engine resolves headers before requests and redirects, and retries a rejected request once when the bearer changes, preserving the active player item across token rotation. diff --git a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift index 1252d836b..70a733130 100644 --- a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift +++ b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift @@ -21,8 +21,10 @@ import Foundation /// The engine owns Range, Host, and HTTP framing. A nil rejected-header dictionary asks for a new /// request. A nonnil dictionary is the actual request headers rejected by one HTTP 401; returning a /// changed Authorization value permits one retry. Other failures and unchanged credentials do not -/// retry. Resolvers may suspend on an actor and should honor task cancellation. The engine bounds -/// each wait and ignores late results after cancellation. Equality compares provider identity. +/// retry. Resolvers may suspend on an actor and should honor task cancellation. They run on one +/// engine-owned serial executor, off Swift's cooperative pool, so they should suspend rather than +/// block. The engine bounds each wait and ignores late results after cancellation. Equality +/// compares provider identity. public final class HTTPRequestAuthorization: Sendable, Equatable { public typealias Resolver = @Sendable (URL, _ rejectedHeaders: [String: String]?) async throws -> [String: String] let resolver: Resolver @@ -57,9 +59,13 @@ public final class HTTPRequestAuthorization: Sendable, Equatable { } } -/// A bridge for the relay's socket worker only. Neither MainActor nor a URLSession delegate queue -/// waits here. A deadline/cancel ends the wait even when the host ignores cancellation; no structured -/// task group waits for an uncooperative child to return. +/// A bridge for the relay's socket worker and the byte-range reader. Neither MainActor nor a +/// URLSession delegate queue waits here. A deadline/cancel ends the wait even when the host ignores +/// cancellation; no structured task group waits for an uncooperative child to return. +/// +/// The waiter can be a cooperative-pool thread: every demuxer open runs inside a `Task.detached`. +/// So the resolver runs on `ResolverExecutor`, never on that pool, or opens filling every pool +/// thread would leave their own resolvers nowhere to start and each would fail at its bound. final class HTTPAuthorizationWait: @unchecked Sendable { private let condition = NSCondition() private var result: Result<[String: String], Error>? @@ -69,12 +75,17 @@ final class HTTPAuthorizationWait: @unchecked Sendable { rejectedHeaders: [String: String]?, timeout: TimeInterval) throws -> [String: String] { condition.lock() if result == nil { - task = Task.detached { [self] in + let operation: @Sendable () async -> Void = { [self] in let answer: Result<[String: String], Error> do { answer = .success(try await authorization.resolver(url, rejectedHeaders)) } catch { answer = .failure(error) } complete(answer) } + if #available(iOS 18, tvOS 18, macOS 15, visionOS 2, *) { + task = Task.detached(executorPreference: ResolverExecutor.shared, operation: operation) + } else { + task = Task.detached(operation: operation) + } } let deadline = Date().addingTimeInterval(timeout) while result == nil { @@ -103,6 +114,20 @@ final class HTTPAuthorizationWait: @unchecked Sendable { } } +/// Runs resolvers, and the default actors they await, off the cooperative pool. A serial queue +/// draws its thread from the overcommit root, so it starts while every pool thread is parked. +/// One process-wide instance: a task holds its preferred executor for its whole life. +@available(iOS 18, tvOS 18, macOS 15, visionOS 2, *) +private final class ResolverExecutor: TaskExecutor { + static let shared = ResolverExecutor() + private let queue = DispatchQueue(label: "com.aetherengine.authorization.resolver") + + func enqueue(_ job: consuming ExecutorJob) { + let job = UnownedJob(job) + queue.async { job.runSynchronously(on: self.asUnownedTaskExecutor()) } + } +} + /// The provider as the byte-range reader sees it: synchronous, bounded, and always asked about the /// source URL. The reader calls it from the demux thread and its probe threads, so every wait ends at /// `timeout` or at `cancel()`, whichever comes first, even when the host never answers. diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index e91de622e..61adeef6f 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -127,6 +127,29 @@ struct RefreshableDirectPlayAuthorizationTests { #expect(!server.requests.contains { $0.range == Self.secondRange }) } + /// Review of PR #12: every open the engine starts runs `Demuxer.open` inside a `Task.detached`, + /// so it parks a cooperative-pool thread while it waits for the resolver. When the resolver + /// needed that same pool, opens that occupied every pool thread left it nowhere to run, and each + /// one failed at its bound instead of starting. Twice the pool's width of parked callers, and a + /// resolver that hops onto an actor, is that state on any machine. + @Test("callers parked on every cooperative thread still get the resolver's answer") + func resolverRunsWhileThePoolIsParked() async throws { + let store = TokenStore() + let authorizer = SourceRequestAuthorizer( + HTTPRequestAuthorization { _, _ in ["Authorization": await store.current()] }, + sourceURL: URL(string: "http://127.0.0.1/movie.mkv")!, timeout: 5) + let callers = ProcessInfo.processInfo.activeProcessorCount * 2 + + let answered = await withTaskGroup(of: Bool.self) { group in + for _ in 0.. String { "Bearer pooled" } +} + /// What the provider was asked. `answer` numbers every call from 1 and records the Authorization of /// each set of rejected headers it was handed. private final class ProviderLog: @unchecked Sendable { diff --git a/docs/api.md b/docs/api.md index 60f48c195..fac6b4240 100644 --- a/docs/api.md +++ b/docs/api.md @@ -146,7 +146,9 @@ ownership, and share refresh work with its API client. Playlist discovery grants authority. Return current credentials without waiting for a proactive refresh while they remain valid; wait for refresh when a credential has expired or was rejected. Never place credentials in URLs. The engine owns Range, routing and HTTP framing headers. Authorization waits are bounded; -stopping the load cancels pending work and ignores late resolver results. +stopping the load cancels pending work and ignores late resolver results. The resolver runs on an +engine-owned serial executor rather than Swift's shared cooperative pool, so engine threads waiting +for its answer cannot starve it; suspend in it rather than block. **Redirect credential scope includes the scheme.** Validate the destination's scheme, host and effective port before obtaining credentials, as well as any session/path restrictions the host From 206f8e66370ee165f341487654ac5d67eabba9a3 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:00:07 -0400 Subject: [PATCH 4/9] fix(avio): latch authorization refusals on backward detour reads detourFetchBlock caught a provider refusal or timeout in its generic catch and returned .failed. serveFromDetour turned that into .miss, so readPersistent fell back to a seek reconnect that asked the provider a second time: a second wait of up to 10 s after an explicit refusal. A 401 on a detour block also became .miss, so the rejected headers never reached the provider and the reconnect sent another request with whatever it answered next. A detour fetch now follows the pump's contract. A provider that refuses or does not answer latches the refusal and fails the read without sending anything. A 401 drops an expired pin as before, then asks the provider with the rejected headers and retries the block once when the Authorization value changes; an unchanged credential or a second 401 latches the 401 and fails the read. The pump and the detour share one helper for that failure. Tests seek more than 8 MiB forward on a 16 MiB source with a 1 MiB window, then read behind it, for a refusal (one provider call, no request), a refreshed 401 (two block requests, stale then fresh), and an unchanged 401 (one block request, read fails). All three failed before this change. Co-Authored-By: Claude Opus 5.5 (1M context) --- Sources/AetherEngine/Demuxer/AVIOReader.swift | 115 ++++++++++++------ ...reshableDirectPlayAuthorizationTests.swift | 112 ++++++++++++++++- 2 files changed, 184 insertions(+), 43 deletions(-) diff --git a/Sources/AetherEngine/Demuxer/AVIOReader.swift b/Sources/AetherEngine/Demuxer/AVIOReader.swift index 2558663b2..971f98a8a 100644 --- a/Sources/AetherEngine/Demuxer/AVIOReader.swift +++ b/Sources/AetherEngine/Demuxer/AVIOReader.swift @@ -1338,7 +1338,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { do { return try authorizer.headers() } catch { - latchProviderRefusal() + latchAuthorizationRefusal() throw openFailureForAuthorization(status: 0) } } @@ -1373,13 +1373,21 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return status == 0 ? .authorizationUnavailable : .httpStatus(status) } - /// The provider refused or did not answer. Latches until a generation delivers again. - private func latchProviderRefusal() { + /// The provider refused or did not answer (status 0), or a 401 stood. Latches until a generation + /// delivers again. + private func latchAuthorizationRefusal(status: Int = 0) { winCond.lock() - authorizationRecovery = .refused(status: 0) + authorizationRecovery = .refused(status: status) winCond.unlock() } + /// The provider has answered, so reconnecting cannot help: the read ends with what it has. + private func failReadForAuthorization(at offset: Int64, status: Int, totalRead: Int) -> Int32 { + EngineLog.emit("[AVIOReader] \(label) authorization refused at offset \(offset) status=\(status); failing the read", category: .demux) + emitNetworkPhase(.exhausted) + return totalRead > 0 ? Int32(totalRead) : -1 + } + /// The streaming GET was answered with a status instead of a body, or the ranged open was /// already refused with one and the unranged GET then delivered nothing either. Either way the /// demuxer would be handed an empty stream (or an error page) and report it as invalid data; @@ -2129,6 +2137,9 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // Hard transport failure: degrade to the OLD single-reconnect behavior. timedReconnect(seek: true, at: curPosition) continue + case .authorizationRefused(let status): + diag.recordDetourFetchAttempt(ms: msSince(detourStart)) + return failReadForAuthorization(at: curPosition, status: status, totalRead: totalRead) } } timedReconnect(seek: true, at: curPosition) @@ -2240,8 +2251,8 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // (#380/#410). detourTrackSequential(at: curPosition, length: n) continue - case .rateLimited, .miss: - break // allowFetch:false never rate-limits; a miss falls through to reconnect + case .rateLimited, .miss, .authorizationRefused: + break // allowFetch:false never fetches; a miss falls through to reconnect } } timedReconnect(seek: true, at: curPosition) @@ -2307,9 +2318,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { timedReconnect(seek: false, at: frontier) continue case .refused(let refusedStatus): - EngineLog.emit("[AVIOReader] \(label) authorization refused at offset \(frontier) status=\(refusedStatus); failing the read", category: .demux) - emitNetworkPhase(.exhausted) // the provider has answered; reconnecting cannot help - return totalRead > 0 ? Int32(totalRead) : -1 + return failReadForAuthorization(at: frontier, status: refusedStatus, totalRead: totalRead) case .notApplicable: break } @@ -2584,8 +2593,14 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// `fetched` says whether the served bytes crossed the network. The callers charge the /// reconnect ladders on it: a resident-block hit is a memcpy out of read-ahead already paid /// for, so it is no more "progress" against a refusing origin than a window serve is (#380). - private enum DetourServe { case served(Int, fetched: Bool); case rateLimited(TimeInterval); case miss } - private enum DetourFetch { case ok(Data); case rateLimited(TimeInterval); case failed } + /// `authorizationRefused` is the latched refusal: the read fails rather than reconnecting. + private enum DetourServe { + case served(Int, fetched: Bool); case rateLimited(TimeInterval); case miss + case authorizationRefused(status: Int) + } + private enum DetourFetch { + case ok(Data); case rateLimited(TimeInterval); case failed; case authorizationRefused(status: Int) + } /// Serve `[offset, offset+maxLen)` (clamped to one 4 MB block) from the detour cache, /// fetching the block over the pooled keep-alive chunkSession on a miss when `allowFetch`. @@ -2621,6 +2636,8 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return .rateLimited(retryAfter) case .failed: return .miss + case .authorizationRefused(let status): + return .authorizationRefused(status: status) } let inBlock = Int(offset - blockStart) @@ -2642,36 +2659,54 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { for: requestURL(), label: "\(label) detour", timeout: budget) defer { OriginRequestBudget.shared.release(ticket) } let rangeEnd = offset + Int64(size) - 1 - var request = URLRequest(url: requestURL()) - request.setValue("bytes=\(offset)-\(rangeEnd)", forHTTPHeaderField: "Range") - // #93/#96: a starved backward-scrub detour fetch must abort fast (the rescue reconnect serves - // instantly), so this path uses the tight interactive budget, not the full chunk timeout. - request.timeoutInterval = budget - do { - let sentHeaders = try applyExtraHeaders(&request) - let (data, response) = try syncRequest(request, headers: sentHeaders, budget: budget) - if let http = response as? HTTPURLResponse { - let status = http.statusCode - if Self.isRateLimitStatus(status) { - let retryAfter = Self.parseRetryAfter(http) - noteOriginRefusal(status: status, retryAfter: retryAfter > 0 ? retryAfter : nil, - respondedBy: http.url) - return .rateLimited(retryAfter) - } - if status != 200 && status != 206 { - if Self.isResolvedExpiryStatus(status) { invalidateResolvedURL() } - return .failed - } - // VOD: 200 at offset > 0 = server ignored Range; silent corruption. Reject. - if status == 200 && offset > 0 && !isLive { - EngineLog.emit("[AVIOReader] detour: server ignored Range (200 for offset \(offset)); rejecting", category: .demux, level: .verbose) - return .failed + // The provider's answer to a 401 on this block: set, the one retry it permits is spent. + var refreshedHeaders: [String: String]? + while true { + var request = URLRequest(url: requestURL()) + request.setValue("bytes=\(offset)-\(rangeEnd)", forHTTPHeaderField: "Range") + // #93/#96: a starved backward-scrub detour fetch must abort fast (the rescue reconnect serves + // instantly), so this path uses the tight interactive budget, not the full chunk timeout. + request.timeoutInterval = budget + guard let sentHeaders = try? applyExtraHeaders(&request, base: refreshedHeaders) else { + // Nothing was sent. A reconnect would only ask the provider again. + latchAuthorizationRefusal() + return .authorizationRefused(status: 0) + } + do { + let (data, response) = try syncRequest(request, headers: sentHeaders, budget: budget) + if let http = response as? HTTPURLResponse { + let status = http.statusCode + if Self.isRateLimitStatus(status) { + let retryAfter = Self.parseRetryAfter(http) + noteOriginRefusal(status: status, retryAfter: retryAfter > 0 ? retryAfter : nil, + respondedBy: http.url) + return .rateLimited(retryAfter) + } + if status != 200 && status != 206 { + if Self.isResolvedExpiryStatus(status) { invalidateResolvedURL() } + // The pump's 401 contract: one retry with a changed credential, then it stands. + if status == 401, let authorizer { + if refreshedHeaders == nil, + let fresh = authorizer.refreshed(rejecting: sentHeaders) { + refreshedHeaders = fresh + continue + } + latchAuthorizationRefusal(status: 401) + return .authorizationRefused(status: 401) + } + return .failed + } + // VOD: 200 at offset > 0 = server ignored Range; silent corruption. Reject. + if status == 200 && offset > 0 && !isLive { + EngineLog.emit("[AVIOReader] detour: server ignored Range (200 for offset \(offset)); rejecting", category: .demux, level: .verbose) + return .failed + } } + addBytesFetched(data.count) + return .ok(data) + } catch { + return .failed } - addBytesFetched(data.count) - return .ok(data) - } catch { - return .failed } } @@ -3580,7 +3615,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { do { sentHeaders = try applyExtraHeaders(&request) } catch { - latchProviderRefusal() + latchAuthorizationRefusal() streamLock.lock() streamEnded = true streamLock.unlock() diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index 61adeef6f..797da5caf 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -127,6 +127,77 @@ struct RefreshableDirectPlayAuthorizationTests { #expect(!server.requests.contains { $0.range == Self.secondRange }) } + // A read behind the window goes to the detour fetch, not the pump. Review of PR #12: the detour + // treated a provider refusal as a transport failure, so the read fell back to a reconnect that + // asked the provider again, and a 401 there never reached the provider as a rejection. + + @Test("a provider refusal on a backward read fails the read without asking again") + func detourRefusalLatches() async throws { + let provider = ProviderLog() + let asksAfterRefusal = Counter() + let server = try Self.origin(total: Self.detourTotal) { _ in true } + defer { server.stop() } + let reader = Self.detourReader(server, authorization: HTTPRequestAuthorization { _, rejected in + if provider.rotated { + asksAfterRefusal.increment() + throw URLError(.userAuthenticationRequired) + } + return provider.answer(rejected: rejected) { "Bearer \($0)" } + }) + defer { reader.markClosed(); reader.close() } + + try await offThread(reader) { try Self.anchorPastTheHead(reader) } + let sentBefore = server.requests.count + provider.rotate() + let read = try await offThread(reader) { Self.readBehindTheWindow(reader) } + + #expect(read.result < 0) + #expect(asksAfterRefusal.value == 1, "the provider was asked \(asksAfterRefusal.value) times") + #expect(server.requests.count == sentBefore, "a request went out after the refusal") + } + + @Test("a 401 on a backward read is retried once with the refreshed credential") + func detourUnauthorizedRetriedWithFreshCredential() async throws { + let provider = ProviderLog() + let server = try Self.origin(total: Self.detourTotal) { + $0.range != Self.detourRange || $0.authorization == "Bearer fresh" + } + defer { server.stop() } + let reader = Self.detourReader(server, authorization: HTTPRequestAuthorization { _, rejected in + provider.answer(rejected: rejected) { _ in rejected == nil ? "Bearer stale" : "Bearer fresh" } + }) + defer { reader.markClosed(); reader.close() } + + try await offThread(reader) { try Self.anchorPastTheHead(reader) } + let sentBefore = server.requests.count + let read = try await offThread(reader) { Self.readBehindTheWindow(reader) } + + #expect(read.result > 0) + let after = server.requests.dropFirst(sentBefore) + #expect(after.map(\.range) == [Self.detourRange, Self.detourRange], "\(after.map(\.range))") + #expect(after.map(\.authorization) == ["Bearer stale", "Bearer fresh"]) + #expect(provider.rejections == ["Bearer stale"]) + } + + @Test("an unchanged credential after a 401 on a backward read fails the read") + func detourUnchangedCredentialFailsTheRead() async throws { + let provider = ProviderLog() + let server = try Self.origin(total: Self.detourTotal) { $0.range != Self.detourRange } + defer { server.stop() } + let reader = Self.detourReader(server, authorization: HTTPRequestAuthorization { _, rejected in + provider.answer(rejected: rejected) { _ in "Bearer stale" } + }) + defer { reader.markClosed(); reader.close() } + + try await offThread(reader) { try Self.anchorPastTheHead(reader) } + let sentBefore = server.requests.count + let read = try await offThread(reader) { Self.readBehindTheWindow(reader) } + + #expect(read.result < 0) + #expect(server.requests.dropFirst(sentBefore).map(\.range) == [Self.detourRange]) + #expect(provider.rejections == ["Bearer stale"]) + } + /// Review of PR #12: every open the engine starts runs `Demuxer.open` inside a `Task.detached`, /// so it parks a cooperative-pool thread while it waits for the resolver. When the resolver /// needed that same pool, opens that occupied every pool thread left it nowhere to run, and each @@ -153,17 +224,18 @@ struct RefreshableDirectPlayAuthorizationTests { // MARK: - Support /// A ranged origin over `total` filler bytes. `accepts` decides per request; a refusal is a 401. - private static func origin(accepts: @escaping @Sendable (ScriptedOriginServer.Recorded) -> Bool) throws + private static func origin(total: Int = total, + accepts: @escaping @Sendable (ScriptedOriginServer.Recorded) -> Bool) throws -> ScriptedOriginServer { try #require(ScriptedOriginServer { request in guard accepts(request) else { return .init(status: 401, declaredLength: 0) } - let (start, end) = Self.bounds(request.range) + let (start, end) = Self.bounds(request.range, total: total) return .init(status: 206, declaredLength: Int64(end - start + 1), contentRange: "bytes \(start)-\(end)/\(total)", bodyBytes: end - start + 1) }) } - private static func bounds(_ range: String?) -> (Int, Int) { + private static func bounds(_ range: String?, total: Int) -> (Int, Int) { let spec = range.map { String($0.dropFirst("bytes=".count)) } ?? "0-" if spec.hasPrefix("-") { return (total - (Int(spec.dropFirst()) ?? 0), total - 1) } let parts = spec.split(separator: "-", omittingEmptySubsequences: false) @@ -181,6 +253,33 @@ struct RefreshableDirectPlayAuthorizationTests { boundedInitialFetch: Int64(openBytes)) } + /// The detour tests need a gap behind the window, so their source is large enough to seek more + /// than 8 MiB past a window held to 1 MiB. + private static let detourTotal = 16 * 1024 * 1024 + /// The 4 MiB detour block that holds byte 9 MiB. + private static let detourRange = "bytes=8388608-12582911" + + private static func detourReader(_ server: ScriptedOriginServer, + authorization: HTTPRequestAuthorization) -> AVIOReader { + AVIOReader(url: URL(string: "http://127.0.0.1:\(server.port)/movie.mkv")!, + requestAuthorization: authorization, boundedInitialFetch: Int64(openBytes), + windowHighWater: 1024 * 1024) + } + + /// Opens, then seeks far enough forward that the pump re-anchors at 12 MiB. + private static func anchorPastTheHead(_ reader: AVIOReader) throws { + try reader.open() + _ = read(reader, count: openBytes) + _ = reader.seek(offset: 12 * 1024 * 1024, whence: SEEK_SET) + _ = read(reader, count: 16 * 1024) + } + + /// A read at 9 MiB, behind the re-anchored window: it goes to the detour fetch. + private static func readBehindTheWindow(_ reader: AVIOReader) -> (bytes: Int, result: Int32) { + _ = reader.seek(offset: 9 * 1024 * 1024, whence: SEEK_SET) + return read(reader, count: 16 * 1024) + } + /// Reads until `count` bytes or the first result that is not a delivery. private static func read(_ reader: AVIOReader, count: Int) -> (bytes: Int, result: Int32) { var buffer = [UInt8](repeating: 0, count: 64 * 1024) @@ -210,6 +309,13 @@ struct RefreshableDirectPlayAuthorizationTests { } } +private final class Counter: @unchecked Sendable { + private let lock = NSLock() + private var count = 0 + var value: Int { lock.withLock { count } } + func increment() { lock.withLock { count += 1 } } +} + /// A host's token store: the resolver awaits it, as a resolver that shares refresh work does. private actor TokenStore { func current() -> String { "Bearer pooled" } From 93ba05d8c774138fd1c44d962650888e1a5fccdd Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:05:19 -0400 Subject: [PATCH 5/9] fix(avio): keep every resolver header off cross-origin targets The byte-range reader asks HTTPRequestAuthorization about the source URL only, then ran its answer through the static-header redirect policy. That policy strips six named credential headers and replays the rest, so a custom credential from the resolver (X-Api-Key, Silo's X-Profile-Id/X-Profile-Token) reached a cross-origin redirect target and every request built against a target pinned from one. The held source connection follows its redirects inline and replayed its whole header set to the next hop, static Authorization included. RedirectHeaderPolicy.Headers now carries two sets per request chain: what a trusted hop gets and what an untrusted one gets. Static headers keep their old split (everything vs. all but the named credentials). A resolver answer is credential-scoped as a whole: a cross-origin hop or pinned target gets the non-credential static headers instead, which is what it would get without a resolver. The reader scopes the set to the request's target, its URLSession delegates and the held connection apply it per hop, and a redirect scrubs every credential-scoped header URLSession carried over. Same-origin requests and the http-to-https upgrade keep the full answer, and readers without a resolver send the same headers as before. The HLS relay is unchanged: it asks the resolver about every hop, so the resolver decides each destination's scope there. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 3 +- Sources/AetherEngine/Demuxer/AVIOReader.swift | 63 ++++++++++-------- .../Demuxer/HeldSourceConnection.swift | 11 +++- .../Demuxer/RedirectHeaderPolicy.swift | 65 ++++++++++++++++--- .../Network/HTTPRequestAuthorization.swift | 6 +- .../Issue377HeldConnectionTests.swift | 10 +-- .../RedirectHeaderPolicyTests.swift | 38 +++++++++++ ...reshableDirectPlayAuthorizationTests.swift | 48 ++++++++++++++ docs/api.md | 8 ++- 9 files changed, 200 insertions(+), 52 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c92b6e4a1..7a094d833 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,12 +23,13 @@ the public-API contract. - `ExternalSubtitleTrack.httpRequestAuthorization` supplies refreshable headers for primary/secondary sidecars and native subtitle stores without changing registered track IDs or rendition mappings. Authorized container decoding retains AVIO streaming and range access. - `HTTPRequestAuthorization.data(from:maximumBytes:)` fetches raw auxiliary resources such as font bundles with a caller-supplied byte limit and a whole-transfer deadline, reusing the relay's redirect, authorization, retry, cancellation and TLS policy. -- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. Resolved credentials follow the static-header redirect policy, so they never reach a cross-origin redirect target. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. Resolvers run on an engine-owned serial executor, off Swift's cooperative pool, so demuxer opens that occupy every pool thread while they wait cannot starve the resolver they wait for. +- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. The resolver is asked about the source only, so every header it returns counts as a credential: none reaches a cross-origin redirect target or a target pinned from one, which get the non-credential static headers instead. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. Resolvers run on an engine-owned serial executor, off Swift's cooperative pool, so demuxer opens that occupy every pool thread while they wait cannot starve the resolver they wait for. - `LoadOptions.httpRequestAuthorization` accepts an async `HTTPRequestAuthorization` resolver for native HLS. The engine resolves headers before requests and redirects, and retries a rejected request once when the bearer changes, preserving the active player item across token rotation. ### Fixed +- `LoadOptions.heldSourceConnection` applies the redirect credential policy to the redirects it follows itself. It replayed every header, `Authorization` and the Emby/Jellyfin tokens included, to a cross-origin hop. - The software video decoder no longer runs more than 16 frame threads. It used one per core, and each frame thread holds back one decoded frame, so a 32-core Mac waited for 31 frames before showing the first one after a load or seek (about 3 s at 10 fps). FFmpeg also warns above 16 threads. Hosts with 16 or fewer cores keep their current thread count. - A paused video no longer starts playing by itself. When the player item died while paused (`failedToPlayToEndTime`), the recovery reload bypassed the pause guard and called `play()` on the fresh item. The reload now keeps a pause made before the item died, whether it came through the engine, AVKit, Control Center or PiP, and mounts the item paused at the same position. - A dead item's recovery no longer restarts the title from where the session was first opened. When AVPlayer refused the recovery item's master (`-11868`), the media fallback reloaded at the first mount's start position, so a title opened from its beginning restarted at 0:00. The fallback now reloads where the refused item was placed. Upstream #621. diff --git a/Sources/AetherEngine/Demuxer/AVIOReader.swift b/Sources/AetherEngine/Demuxer/AVIOReader.swift index 971f98a8a..723e247c9 100644 --- a/Sources/AetherEngine/Demuxer/AVIOReader.swift +++ b/Sources/AetherEngine/Demuxer/AVIOReader.swift @@ -1091,22 +1091,29 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// both `Authorization` and `X-Emby-Token`. One policy, applied where the request is built, so /// a pin cannot outflank it. /// - /// With a provider, `base` is its answer for the source URL, asked once per request: the same - /// policy then decides where that answer may go, exactly as it does for static headers. A throw - /// means the provider refused, timed out or the reader closed, and the request must not be sent. - private func headers(for target: URL?, base: [String: String]? = nil) throws -> [String: String] { - let source = try base ?? authorizer?.headers() ?? extraHeaders - return RedirectHeaderPolicy.headersToReplay( - extraHeaders: source, originalURL: url, redirectURL: target ?? url) + /// With a provider, `base` is its answer for the source URL, asked once per request. The + /// provider is never asked about the target, so every header of that answer counts as a + /// credential: a target the source's credentials may not reach gets the static headers that are + /// not credentials instead. A throw means the provider refused, timed out or the reader closed, + /// and the request must not be sent. + private func headers(for target: URL?, base: [String: String]? = nil) throws -> RedirectHeaderPolicy.Headers { + let headers: RedirectHeaderPolicy.Headers + if let answer = try base ?? authorizer?.headers() { + headers = .init(authorized: answer, static: extraHeaders) + } else { + headers = .init(static: extraHeaders) + } + return headers.scoped(source: url, target: target ?? url) } /// Sets the request's headers and returns them, so the delegate replays this request's own - /// headers on a redirect instead of asking the provider a second time. + /// headers on a redirect instead of asking the provider a second time. `credentialed` is what + /// this request carries. @discardableResult private func applyExtraHeaders(_ request: inout URLRequest, - base: [String: String]? = nil) throws -> [String: String] { + base: [String: String]? = nil) throws -> RedirectHeaderPolicy.Headers { let headers = try headers(for: request.url, base: base) - for (name, value) in headers { + for (name, value) in headers.credentialed { request.setValue(value, forHTTPHeaderField: name) } return headers @@ -2687,7 +2694,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { // The pump's 401 contract: one retry with a changed credential, then it stands. if status == 401, let authorizer { if refreshedHeaders == nil, - let fresh = authorizer.refreshed(rejecting: sentHeaders) { + let fresh = authorizer.refreshed(rejecting: sentHeaders.credentialed) { refreshedHeaders = fresh continue } @@ -3126,7 +3133,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { } request.timeoutInterval = 0 // long-lived; stalls handled by the reader // Asked before the origin slot is taken, so a slow provider never holds one. - let sentHeaders: [String: String] + let sentHeaders: RedirectHeaderPolicy.Headers do { sentHeaders = try applyExtraHeaders(&request, base: staged) } catch { @@ -3176,7 +3183,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return } activeTransfer = transfer - connSentHeaders = sentHeaders + connSentHeaders = sentHeaders.credentialed winCond.unlock() transfer.startTransfer() @@ -3611,7 +3618,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { private func streamDownloadSync() { var request = URLRequest(url: url) request.timeoutInterval = 0 // No timeout for live streams - let sentHeaders: [String: String] + let sentHeaders: RedirectHeaderPolicy.Headers do { sentHeaders = try applyExtraHeaders(&request) } catch { @@ -4098,7 +4105,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { if let http = response as? HTTPURLResponse { let status = http.statusCode if status == 401, refreshedHeaders == nil, - let fresh = authorizer?.refreshed(rejecting: sentHeaders) { + let fresh = authorizer?.refreshed(rejecting: sentHeaders.credentialed) { return fetchChunkAttempt(from: offset, size: size, forceSource: forceSource, refreshedHeaders: fresh) } @@ -4283,7 +4290,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { } /// `headers` are the ones `applyExtraHeaders` set on `request`, replayed on a redirect. - private func syncRequest(_ request: URLRequest, headers: [String: String], + private func syncRequest(_ request: URLRequest, headers: RedirectHeaderPolicy.Headers, budget: TimeInterval = 35) throws -> (Data, URLResponse) { // #377: every short fetch the reader makes (detour blocks, size probes, HEAD) funnels // through here, so this is the one place that has to take an origin slot for all of them. @@ -4416,7 +4423,7 @@ final class DetourBlockCache: @unchecked Sendable { private func redirectPreservingHeaders( task: URLSessionTask, newRequest request: URLRequest, - extraHeaders: [String: String] + extraHeaders: RedirectHeaderPolicy.Headers ) -> URLRequest { // #388: this is the moment the request the reader budgeted for stops being answered by the // origin it was budgeted against. Every fetch the reader makes passes through here, so it is @@ -4429,7 +4436,7 @@ private func redirectPreservingHeaders( request, originalURL: task.originalRequest?.url, originalRange: task.originalRequest?.value(forHTTPHeaderField: "Range"), - extraHeaders: extraHeaders) + headers: extraHeaders) } // MARK: - Persistent Read Delegate @@ -4579,7 +4586,7 @@ extension URLSessionDataTask: PersistentTransfer { private final class PersistentReadDelegate: NSObject, URLSessionDataDelegate, @unchecked Sendable { weak var reader: AVIOReader? let generation: Int - let extraHeaders: [String: String] + let extraHeaders: RedirectHeaderPolicy.Headers /// #377: the origin slot this connection occupies, held here because the delegate's lifetime /// IS the task's. Seven paths in the reader clear `activeTransfer` and only one of them is the /// task ending, so a ticket released alongside `activeTransfer` would leak on the other six. @@ -4589,7 +4596,7 @@ private final class PersistentReadDelegate: NSObject, URLSessionDataDelegate, @u /// The URL this connection was opened against, for the one-per-origin transport line. private let originURL: URL - init(reader: AVIOReader, generation: Int, extraHeaders: [String: String], + init(reader: AVIOReader, generation: Int, extraHeaders: RedirectHeaderPolicy.Headers, ticket: OriginRequestBudget.Ticket?, originURL: URL) { self.reader = reader self.generation = generation @@ -4696,7 +4703,7 @@ private final class PersistentReadDelegate: NSObject, URLSessionDataDelegate, @u /// dispatch_data is released per delivery. @unchecked Sendable: ownership /// via semaphore ensures no concurrent access to mutable fields. private final class ChunkFetchDelegate: NSObject, URLSessionDataDelegate, @unchecked Sendable { - let extraHeaders: [String: String] + let extraHeaders: RedirectHeaderPolicy.Headers /// Most this fetch will ever buffer, nil when the request is open-ended (#255). Derived from /// the request, never from the response: a declared length is the origin's claim about the /// whole source, not about this body. @@ -4710,7 +4717,7 @@ private final class ChunkFetchDelegate: NSObject, URLSessionDataDelegate, @unche var onCompletion: (() -> Void)? var onResolved: ((URL) -> Void)? - init(extraHeaders: [String: String], bodyLimit: Int?) { + init(extraHeaders: RedirectHeaderPolicy.Headers, bodyLimit: Int?) { self.extraHeaders = extraHeaders self.bodyLimit = bodyLimit } @@ -4818,10 +4825,10 @@ private final class StreamingDelegate: NSObject, URLSessionDataDelegate { /// Re-applied across cross-host redirects like every other delegate in this file; /// IPTV origins routinely 302 twice (portal -> panel -> archive host) and the final /// host must still see the caller's User-Agent / auth headers. - let extraHeaders: [String: String] + let extraHeaders: RedirectHeaderPolicy.Headers init( - extraHeaders: [String: String] = [:], + extraHeaders: RedirectHeaderPolicy.Headers = .init(static: [:]), onResponse: (@Sendable (URLResponse) -> Void)? = nil, onRefused: (@Sendable (Int, URL?) -> Void)? = nil, onData: @escaping @Sendable (Data) -> Void, @@ -4894,12 +4901,12 @@ private final class StreamingDelegate: NSObject, URLSessionDataDelegate { /// captures total from Content-Range, cancels before the body streams. /// @unchecked Sendable: single-use per probe, semaphore ownership prevents concurrency. private final class ProbeDelegate: NSObject, URLSessionDataDelegate, @unchecked Sendable { - let extraHeaders: [String: String] + let extraHeaders: RedirectHeaderPolicy.Headers var totalSize: Int64? var onCompletion: (() -> Void)? var onResolved: ((URL) -> Void)? - init(extraHeaders: [String: String]) { + init(extraHeaders: RedirectHeaderPolicy.Headers) { self.extraHeaders = extraHeaders } @@ -5046,7 +5053,7 @@ private final class TailPrefetchDelegate: NSObject, URLSessionDataDelegate, @unc } private let expectedLength: Int - private let extraHeaders: [String: String] + private let extraHeaders: RedirectHeaderPolicy.Headers private var buffer = Data() private var spanStart: Int64? private var rejection: (String, Verdict)? @@ -5055,7 +5062,7 @@ private final class TailPrefetchDelegate: NSObject, URLSessionDataDelegate, @unc /// silent failure would be a caller waiting out its whole budget for bytes that are never coming. var onOutcome: ((Outcome) -> Void)? - init(expectedLength: Int, extraHeaders: [String: String]) { + init(expectedLength: Int, extraHeaders: RedirectHeaderPolicy.Headers) { self.expectedLength = expectedLength self.extraHeaders = extraHeaders } diff --git a/Sources/AetherEngine/Demuxer/HeldSourceConnection.swift b/Sources/AetherEngine/Demuxer/HeldSourceConnection.swift index b24fbf48c..ba6341511 100644 --- a/Sources/AetherEngine/Demuxer/HeldSourceConnection.swift +++ b/Sources/AetherEngine/Demuxer/HeldSourceConnection.swift @@ -83,7 +83,10 @@ final class HeldSourceConnection: @unchecked Sendable { } private weak var delegate: HeldSourceConnectionDelegate? - private let extraHeaders: [String: String] + /// Split by `RedirectHeaderPolicy`: a hop followed inline gets what the policy allows it from + /// `requestedURL`, never the whole set. + private let extraHeaders: RedirectHeaderPolicy.Headers + private let requestedURL: URL private let offset: Int64 private let userAgent: String? private let queue: DispatchQueue @@ -106,13 +109,14 @@ final class HeldSourceConnection: @unchecked Sendable { init(url: URL, offset: Int64, - extraHeaders: [String: String], + extraHeaders: RedirectHeaderPolicy.Headers, userAgent: String?, label: String, generation: Int = 0, ticket: OriginRequestBudget.Ticket? = nil, delegate: HeldSourceConnectionDelegate) { self.respondedBy = url + self.requestedURL = url self.offset = offset self.extraHeaders = extraHeaders self.userAgent = userAgent @@ -237,7 +241,8 @@ final class HeldSourceConnection: @unchecked Sendable { if secure { task.startSecureConnection() } try write(Self.requestBytes(target: target, host: host, port: port, secure: secure, - offset: offset, extraHeaders: extraHeaders, + offset: offset, + extraHeaders: extraHeaders.toReplay(from: requestedURL, to: target), userAgent: userAgent)) return try readHead() } diff --git a/Sources/AetherEngine/Demuxer/RedirectHeaderPolicy.swift b/Sources/AetherEngine/Demuxer/RedirectHeaderPolicy.swift index 185326261..d0e9b465b 100644 --- a/Sources/AetherEngine/Demuxer/RedirectHeaderPolicy.swift +++ b/Sources/AetherEngine/Demuxer/RedirectHeaderPolicy.swift @@ -7,6 +7,10 @@ import Foundation /// the request outright when the target authenticates via the URL itself and rejects /// conflicting auth mechanisms with 400. Non-credential headers are replayed /// unconditionally so header-dependent proxies keep working (#8). +/// +/// Headers from `HTTPRequestAuthorization` are different: the engine asks the provider about the +/// source only, so it cannot tell which of the answer's headers are credentials. Every one of them +/// is treated as one, and an untrusted destination gets what it would without a provider. enum RedirectHeaderPolicy { private static let credentialHeaders: Set = [ "authorization", @@ -17,15 +21,47 @@ enum RedirectHeaderPolicy { "x-mediabrowser-token", ] + /// The headers one request chain may carry, split by where they may go. `credentialed` reaches + /// only a hop the chain's origin may share credentials with; every other hop gets `anonymous`. + struct Headers: Sendable, Equatable { + let credentialed: [String: String] + let anonymous: [String: String] + + private init(credentialed: [String: String], anonymous: [String: String]) { + self.credentialed = credentialed + self.anonymous = anonymous + } + + /// Static headers: all of them to a trusted hop, all but the named credentials elsewhere. + init(static headers: [String: String]) { + self.init(credentialed: headers, anonymous: withoutCredentials(headers)) + } + + /// A provider's answer for the source. None of it goes to an untrusted hop, which gets the + /// static headers that are not credentials instead. + init(authorized answer: [String: String], static headers: [String: String]) { + self.init(credentialed: answer, anonymous: withoutCredentials(headers)) + } + + func toReplay(from original: URL?, to destination: URL?) -> [String: String] { + credentialsAllowed(from: original, to: destination) ? credentialed : anonymous + } + + /// These headers for a request built against `target` on behalf of `source`. A target the + /// source's credentials may not reach, such as one pinned from a cross-origin redirect, + /// carries `anonymous` on its own redirects too. + func scoped(source: URL, target: URL) -> Headers { + credentialsAllowed(from: source, to: target) + ? self : Headers(credentialed: anonymous, anonymous: anonymous) + } + } + static func headersToReplay( extraHeaders: [String: String], originalURL: URL?, redirectURL: URL? ) -> [String: String] { - if credentialsAllowed(from: originalURL, to: redirectURL) { - return extraHeaders - } - return withoutCredentials(extraHeaders) + Headers(static: extraHeaders).toReplay(from: originalURL, to: redirectURL) } static func withoutCredentials(_ headers: [String: String]) -> [String: String] { @@ -42,19 +78,30 @@ enum RedirectHeaderPolicy { originalURL: URL?, originalRange: String?, extraHeaders: [String: String] + ) -> URLRequest { + redirectRequest(request, originalURL: originalURL, originalRange: originalRange, + headers: Headers(static: extraHeaders)) + } + + /// The same, for a chain whose headers are already split. An untrusted hop loses every + /// `credentialed` header URLSession carried over, not only the named credentials. + static func redirectRequest( + _ request: URLRequest, + originalURL: URL?, + originalRange: String?, + headers: Headers ) -> URLRequest { var updated = request if let originalRange { updated.setValue(originalRange, forHTTPHeaderField: "Range") } - if !credentialsAllowed(from: originalURL, to: request.url) { - for name in credentialHeaders { + let trusted = credentialsAllowed(from: originalURL, to: request.url) + if !trusted { + for name in credentialHeaders.union(headers.credentialed.keys) { updated.setValue(nil, forHTTPHeaderField: name) } } - let replayable = headersToReplay( - extraHeaders: extraHeaders, originalURL: originalURL, redirectURL: request.url) - for (name, value) in replayable { + for (name, value) in trusted ? headers.credentialed : headers.anonymous { updated.setValue(value, forHTTPHeaderField: name) } return updated diff --git a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift index 70a733130..6f92e6d0e 100644 --- a/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift +++ b/Sources/AetherEngine/Network/HTTPRequestAuthorization.swift @@ -8,9 +8,9 @@ import Foundation /// transport policy. Live ingest still uses its static headers. /// /// Direct media (the engine's own byte-range reader) asks for the source URL the host loaded before -/// every request it builds: each range, reconnect, probe and seek. The answer then follows the -/// static-header redirect policy, so credentials reach only the source's origin (and an http-to-https -/// upgrade of it), never a cross-origin redirect target or a target pinned from one. +/// every request it builds: each range, reconnect, probe and seek. Because it is asked about nothing +/// else, every header of the answer counts as a credential and reaches only the source's origin (and +/// an http-to-https upgrade of it), never a cross-origin redirect target or a target pinned from one. /// The resolver must independently validate every URL, including redirects and playlist-discovered /// origins. Discovery grants no credential authority. Credentials must never be placed in URLs. /// Include scheme, host and effective port in that scope. Redirects are authorized afresh. Throw to diff --git a/Tests/AetherEngineTests/Issue377HeldConnectionTests.swift b/Tests/AetherEngineTests/Issue377HeldConnectionTests.swift index 4a9c63db9..7f9fddcd3 100644 --- a/Tests/AetherEngineTests/Issue377HeldConnectionTests.swift +++ b/Tests/AetherEngineTests/Issue377HeldConnectionTests.swift @@ -185,7 +185,7 @@ struct Issue377HeldConnectionTests { let url = try #require(URL(string: "http://127.0.0.1:\(origin.port)/media.mkv")) let delegate = RecordingHeldDelegate(defaultBudget: 128 * 1024) - let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: [:], + let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: .init(static: [:]), userAgent: "AetherEngine/test", label: "test", delegate: delegate) connection.start() @@ -208,7 +208,7 @@ struct Issue377HeldConnectionTests { // Two pulls, then the answer a paused viewer produces. let delegate = RecordingHeldDelegate(budgets: [32 * 1024, 32 * 1024, 0]) - let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: [:], + let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: .init(static: [:]), userAgent: nil, label: "test", delegate: delegate) connection.start() #expect(delegate.waitForEnd()) @@ -233,7 +233,7 @@ struct Issue377HeldConnectionTests { let url = try #require(URL(string: "http://127.0.0.1:\(origin.port)/media.mkv")) let delegate = RecordingHeldDelegate() - let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: [:], + let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: .init(static: [:]), userAgent: nil, label: "test", delegate: delegate) connection.start() #expect(delegate.waitForEnd()) @@ -253,7 +253,7 @@ struct Issue377HeldConnectionTests { let url = try #require(URL(string: "http://127.0.0.1:\(origin.port)/source.mkv")) let delegate = RecordingHeldDelegate(defaultBudget: 64 * 1024) - let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: [:], + let connection = HeldSourceConnection(url: url, offset: 0, extraHeaders: .init(static: [:]), userAgent: nil, label: "test", delegate: delegate) connection.start() #expect(delegate.waitForEnd()) @@ -273,7 +273,7 @@ struct Issue377HeldConnectionTests { let offset: Int64 = 1_048_576 let delegate = RecordingHeldDelegate(defaultBudget: 256 * 1024) - let connection = HeldSourceConnection(url: url, offset: offset, extraHeaders: [:], + let connection = HeldSourceConnection(url: url, offset: offset, extraHeaders: .init(static: [:]), userAgent: nil, label: "test", delegate: delegate) connection.start() #expect(delegate.waitForEnd()) diff --git a/Tests/AetherEngineTests/RedirectHeaderPolicyTests.swift b/Tests/AetherEngineTests/RedirectHeaderPolicyTests.swift index 41d8a6fbb..89943d550 100644 --- a/Tests/AetherEngineTests/RedirectHeaderPolicyTests.swift +++ b/Tests/AetherEngineTests/RedirectHeaderPolicyTests.swift @@ -124,6 +124,44 @@ struct RedirectHeaderPolicyTests { #expect(out["X-Custom"] == "1") } + // MARK: Provider answers + + // The provider is asked about the source only, so any header it returns may be a credential. + + private let answer = ["Authorization": "Bearer abc", "X-Profile-Token": "p1", "Referer": "provider"] + private let staticHeaders = ["Referer": "app", "X-Emby-Token": "static"] + + @Test("A cross-origin hop gets no provider header, only the non-credential static ones") + func providerHeadersStayOffCrossOriginHops() { + let headers = RedirectHeaderPolicy.Headers(authorized: answer, static: staticHeaders) + let source = url("https://media.example/x") + #expect(headers.toReplay(from: source, to: url("https://cdn.example/x")) == ["Referer": "app"]) + #expect(headers.toReplay(from: source, to: url("https://media.example/y")) == answer) + } + + @Test("A target pinned from a cross-origin redirect carries no provider header on its own hops") + func pinnedTargetCarriesNoProviderHeaders() { + let pinned = url("https://cdn.example/x") + let headers = RedirectHeaderPolicy.Headers(authorized: answer, static: staticHeaders) + .scoped(source: url("https://media.example/x"), target: pinned) + #expect(headers.credentialed == ["Referer": "app"]) + #expect(headers.toReplay(from: pinned, to: url("https://cdn.example/y")) == ["Referer": "app"]) + } + + @Test("A custom provider header URLSession carried over is removed cross-host") + func carriedOverProviderHeaderRemovedCrossHost() { + var carried = URLRequest(url: url("https://cdn.example/o.mp4")) + carried.setValue("p1", forHTTPHeaderField: "X-Profile-Token") + carried.setValue("provider", forHTTPHeaderField: "Referer") + let out = RedirectHeaderPolicy.redirectRequest( + carried, + originalURL: url("https://media.example/x"), + originalRange: nil, + headers: .init(authorized: answer, static: staticHeaders)) + #expect(out.value(forHTTPHeaderField: "X-Profile-Token") == nil) + #expect(out.value(forHTTPHeaderField: "Referer") == "app") + } + // MARK: Request-level sanitization @Test("Credential carried over by URLSession is removed cross-host") diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index 797da5caf..7aa159df0 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -127,6 +127,54 @@ struct RefreshableDirectPlayAuthorizationTests { #expect(!server.requests.contains { $0.range == Self.secondRange }) } + /// Review of PR #12: the provider is asked only about the source, so the engine cannot know which + /// of its headers are credentials. The static policy stripped six named ones, and a custom one + /// (`X-Profile-Token`, `X-Api-Key`) reached the redirect target and every request built against + /// the target the session pinned from it. The held connection follows its redirects inline and + /// replayed every header to the next hop, static credentials included. + @Test("no header the provider returns reaches a cross-origin target, followed or pinned", + arguments: [false, true]) + func providerHeadersStayOnTheSourceOrigin(heldConnection: Bool) async throws { + let fileSize: Int64 = 64 * 1024 * 1024 + let firstRange = 256 * 1024 + let cdn = try #require(ThrottledOriginServer(totalSize: fileSize)) + defer { cdn.stop() } + let cdnPort = cdn.port + let redirecting = ThrottledOriginServer(totalSize: fileSize, respond: { _, _, _ in + .redirect(to: "http://127.0.0.1:\(cdnPort)/cdn/movie.bin") + }) + let source = try #require(redirecting) + defer { source.stop() } + let reader = AVIOReader( + url: URL(string: "http://127.0.0.1:\(source.port)/movie.bin")!, + extraHeaders: ["Referer": "https://app.example", "X-Emby-Token": "STATIC"], + requestAuthorization: HTTPRequestAuthorization { _, _ in + ["Authorization": "Bearer SOURCE-ONLY", "X-Profile-Token": "SOURCE-ONLY", + "Referer": "https://provider.example"] + }, + boundedInitialFetch: Int64(firstRange), heldConnection: heldConnection) + defer { reader.markClosed(); reader.close() } + + // Past the bounded first range, so the pump also builds a request against the pinned target. + let want = firstRange + 128 * 1024 + let read = try await offThread(reader) { + try reader.open() + return Self.read(reader, count: want) + } + + #expect(read.bytes == want) + let atTarget = cdn.requestHeaders + // The held connection asks for one open-ended range, so it reaches the target only by its + // inline hop; the URLSession pump also builds a request against the pinned target. + #expect(atTarget.count >= (heldConnection ? 1 : 2), "\(atTarget)") + for name in ["authorization", "x-profile-token", "x-emby-token"] { + #expect(atTarget.allSatisfy { $0[name] == nil }, "\(name) reached the target: \(atTarget)") + } + // What a target gets without a provider: the static headers that are not credentials. + #expect(atTarget.allSatisfy { $0["referer"] == "https://app.example" }, "\(atTarget)") + #expect(source.requestHeaders.allSatisfy { $0["x-profile-token"] == "SOURCE-ONLY" }) + } + // A read behind the window goes to the detour fetch, not the pump. Review of PR #12: the detour // treated a provider refusal as a transport failure, so the read fell back to a reconnect that // asked the provider again, and a 401 there never reached the provider as a rejection. diff --git a/docs/api.md b/docs/api.md index fac6b4240..e4f022ad5 100644 --- a/docs/api.md +++ b/docs/api.md @@ -177,9 +177,11 @@ E-AC-3, PCM): the resolver authorizes its probe, then AVPlayer plays it with `ht On direct play the reader asks the resolver for the source URL the host loaded before every request it builds: each range, reconnect, seek, size probe and tail fetch. The answer replaces -`httpHeaders` and then follows the static-header redirect policy, so credentials reach only the -source's origin (or an http-to-https upgrade of it), never a cross-origin redirect target or a -target pinned from one. The resolver is not asked about those destinations. After a 401 the +`httpHeaders` on requests to the source's origin (or an http-to-https upgrade of it). The resolver +is not asked about any other destination, so the engine treats every header it returns as a +credential, custom ones such as `X-Api-Key` included: a cross-origin redirect target, or a target +pinned from one, receives none of them. Such a target gets `httpHeaders` without the credential +headers named above, as it would without a resolver. After a 401 the resolver receives the headers that request carried; a changed `Authorization` value retries the request once at the same byte offset. Unchanged credentials, a second 401, or a resolver that throws or does not answer within its bound end the read instead of running the reconnect ladder, and fail From 276df45a7c438bcd1487428e4b90f302ea610e3a Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:06:07 -0400 Subject: [PATCH 6/9] docs(api): separate pre-send and post-401 direct-play auth failures The direct-play paragraph listed unchanged credentials, a second 401 and a resolver failure together and said they "fail an open before anything is sent". Only the resolver failure happens before a request: the 401 cases end the read after the origin answered, and fail an open with that status rather than as authorizationUnavailable. The paragraph now lists the two groups apart, with what a load reports for each (.sourceOpenFailed for authorizationUnavailable, .sourceRefused with underlyingCode 401 after a 401), and the CHANGELOG entry draws the same line. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- docs/api.md | 12 +++++++++--- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a094d833..af7fa1beb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,7 @@ the public-API contract. - `ExternalSubtitleTrack.httpRequestAuthorization` supplies refreshable headers for primary/secondary sidecars and native subtitle stores without changing registered track IDs or rendition mappings. Authorized container decoding retains AVIO streaming and range access. - `HTTPRequestAuthorization.data(from:maximumBytes:)` fetches raw auxiliary resources such as font bundles with a caller-supplied byte limit and a whole-transfer deadline, reusing the relay's redirect, authorization, retry, cancellation and TLS policy. -- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`; unchanged credentials, a second 401, or a resolver that throws or exceeds its 10 s bound fail the read without running the reconnect ladder. The resolver is asked about the source only, so every header it returns counts as a credential: none reaches a cross-origin redirect target or a target pinned from one, which get the non-credential static headers instead. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. Resolvers run on an engine-owned serial executor, off Swift's cooperative pool, so demuxer opens that occupy every pool thread while they wait cannot starve the resolver they wait for. +- `LoadOptions.httpRequestAuthorization` now covers direct play. The byte-range reader asks the resolver for the source URL before every range, reconnect, probe and seek, so a rotated access token reaches the next request instead of the session sending the headers it opened with until the host reloads the player. A 401 retries once at the same byte offset when the resolver returns a changed `Authorization`. A resolver that throws or exceeds its 10 s bound fails the request before it is sent, so an open fails as `authorizationUnavailable`. Unchanged credentials, a second 401, or a failed refresh end the read after the 401, and fail an open with that status. Neither runs the reconnect ladder. The resolver is asked about the source only, so every header it returns counts as a credential: none reaches a cross-origin redirect target or a target pinned from one, which get the non-credential static headers instead. Live ingest, remote disc images and audio-only sources that AVPlayer decodes natively keep static headers. Resolvers run on an engine-owned serial executor, off Swift's cooperative pool, so demuxer opens that occupy every pool thread while they wait cannot starve the resolver they wait for. - `LoadOptions.httpRequestAuthorization` accepts an async `HTTPRequestAuthorization` resolver for native HLS. The engine resolves headers before requests and redirects, and retries a rejected request once when the bearer changes, preserving the active player item across token rotation. diff --git a/docs/api.md b/docs/api.md index e4f022ad5..635a8dd84 100644 --- a/docs/api.md +++ b/docs/api.md @@ -183,9 +183,15 @@ credential, custom ones such as `X-Api-Key` included: a cross-origin redirect ta pinned from one, receives none of them. Such a target gets `httpHeaders` without the credential headers named above, as it would without a resolver. After a 401 the resolver receives the headers that request carried; a changed `Authorization` value retries the -request once at the same byte offset. Unchanged credentials, a second 401, or a resolver that throws -or does not answer within its bound end the read instead of running the reconnect ladder, and fail -an open before anything is sent. A rotated token therefore needs no player reload. +request once at the same byte offset. A rotated token therefore needs no player reload. Failures +fall into two groups, and neither runs the reconnect ladder: + +- Before a request is sent: the resolver throws or does not answer within its bound. Nothing goes + to the origin. An open fails with `AVIOReaderError.authorizationUnavailable`, which a load reports + as `.sourceOpenFailed`; a read in progress ends. +- After a 401: the credential is unchanged, the retry is refused again, or the resolver throws or + times out while answering the 401. A read in progress ends; an open fails with the 401, which a + load reports as `.sourceRefused` with `underlyingCode` 401. For external subtitles, set `ExternalSubtitleTrack.httpRequestAuthorization` on each registered track. This is independent of the media provider, so the host can restrict subtitle credentials From 2df23a6a6f5291a57ba5b6c748ca5f1ff36e3313 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:18:01 -0400 Subject: [PATCH 7/9] fix(avio): honour a latched authorization refusal before a detour fetch 206f8e66 made a backward detour read latch an authorization refusal, but nothing read the latch on the way back in. A failed read leaves the cursor where it was, so the demuxer's retry took the backward branch to the detour fetch again: the provider was asked a second time, and after a 401 the block was requested again, with no generation delivery in between. The pump path did not have this gap, because its retry reaches recoverAuthorization, which returns the latch at once. detourFetchBlock now returns the latched refusal before it takes an origin slot, asks the provider or builds a request. A resident detour block is still served from memory. latchedAuthorizationRefusal() also replaces the open's inline read of the same latch. The refusal and unchanged-credential detour tests now read the same backward offset twice. Before this change the second read asked the provider again, and after a 401 sent a second block request. Co-Authored-By: Claude Opus 5.5 (1M context) --- Sources/AetherEngine/Demuxer/AVIOReader.swift | 18 ++++++++++++++---- ...freshableDirectPlayAuthorizationTests.swift | 10 ++++++++-- 2 files changed, 22 insertions(+), 6 deletions(-) diff --git a/Sources/AetherEngine/Demuxer/AVIOReader.swift b/Sources/AetherEngine/Demuxer/AVIOReader.swift index 723e247c9..d13b342d5 100644 --- a/Sources/AetherEngine/Demuxer/AVIOReader.swift +++ b/Sources/AetherEngine/Demuxer/AVIOReader.swift @@ -1388,6 +1388,14 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { winCond.unlock() } + /// The status of the latched refusal, nil while none stands. + private func latchedAuthorizationRefusal() -> Int? { + winCond.lock() + defer { winCond.unlock() } + if case .refused(let status) = authorizationRecovery { return status } + return nil + } + /// The provider has answered, so reconnecting cannot help: the read ends with what it has. private func failReadForAuthorization(at offset: Int64, status: Int, totalRead: Int) -> Int32 { EngineLog.emit("[AVIOReader] \(label) authorization refused at offset \(offset) status=\(status); failing the read", category: .demux) @@ -1439,10 +1447,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { close() throw AVIOReaderError.transportSecurityFailed(code: tlsCode) } - winCond.lock() - let authorization = authorizationRecovery - winCond.unlock() - if case .refused(let refusedStatus) = authorization, status == 0 { + if status == 0, let refusedStatus = latchedAuthorizationRefusal() { throw openFailureForAuthorization(status: refusedStatus) } guard status != 0 else { return } @@ -2661,6 +2666,11 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { /// Single Range fetch for a detour block over the pooled chunkSession. Surfaces rate limiting with /// its Retry-After so the caller can back off in place rather than churn the connection (#71). private func detourFetchBlock(from offset: Int64, size: Int) -> DetourFetch { + // A failed read leaves the cursor where it was, so the retry lands here again. Until a + // generation delivers, the latched refusal stands: no provider call, no request. + if let status = latchedAuthorizationRefusal() { + return .authorizationRefused(status: status) + } let budget = Self.effectiveDetourBudget(chunkRequestTimeout: chunkRequestTimeout) let ticket = OriginRequestBudget.shared.acquire( for: requestURL(), label: "\(label) detour", timeout: budget) diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index 7aa159df0..4c2f68637 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -179,7 +179,7 @@ struct RefreshableDirectPlayAuthorizationTests { // treated a provider refusal as a transport failure, so the read fell back to a reconnect that // asked the provider again, and a 401 there never reached the provider as a rejection. - @Test("a provider refusal on a backward read fails the read without asking again") + @Test("a provider refusal on a backward read latches: a repeat read neither asks nor sends") func detourRefusalLatches() async throws { let provider = ProviderLog() let asksAfterRefusal = Counter() @@ -198,8 +198,12 @@ struct RefreshableDirectPlayAuthorizationTests { let sentBefore = server.requests.count provider.rotate() let read = try await offThread(reader) { Self.readBehindTheWindow(reader) } + // A failed read leaves the cursor where it was, so the demuxer's retry lands on the same + // backward offset. Nothing has delivered since, so the refusal must still stand. + let repeated = try await offThread(reader) { Self.readBehindTheWindow(reader) } #expect(read.result < 0) + #expect(repeated.result < 0) #expect(asksAfterRefusal.value == 1, "the provider was asked \(asksAfterRefusal.value) times") #expect(server.requests.count == sentBefore, "a request went out after the refusal") } @@ -227,7 +231,7 @@ struct RefreshableDirectPlayAuthorizationTests { #expect(provider.rejections == ["Bearer stale"]) } - @Test("an unchanged credential after a 401 on a backward read fails the read") + @Test("an unchanged credential after a 401 on a backward read fails it, and a repeat read too") func detourUnchangedCredentialFailsTheRead() async throws { let provider = ProviderLog() let server = try Self.origin(total: Self.detourTotal) { $0.range != Self.detourRange } @@ -240,8 +244,10 @@ struct RefreshableDirectPlayAuthorizationTests { try await offThread(reader) { try Self.anchorPastTheHead(reader) } let sentBefore = server.requests.count let read = try await offThread(reader) { Self.readBehindTheWindow(reader) } + let repeated = try await offThread(reader) { Self.readBehindTheWindow(reader) } #expect(read.result < 0) + #expect(repeated.result < 0) #expect(server.requests.dropFirst(sentBefore).map(\.range) == [Self.detourRange]) #expect(provider.rejections == ["Bearer stale"]) } From ba4728d4b2fc53a11f367f78c70a5223d01a9920 Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:23:39 -0400 Subject: [PATCH 8/9] test(avio): stall the anchored pump in the detour authorization tests The detour tests re-anchored the pump at 12 MiB with a 1 MiB window and let it run. Under the parallel authorization group the pump's backpressure refill (bytes=14417920-) sometimes landed while a test was measuring: it showed up as an extra request, and its delivery lifted the detour's latched refusal, as a delivery is meant to, so the repeat read legitimately asked the provider again. The detour tests failed under group load and passed alone. A backward read takes the detour only while the pump is connected, so the anchored range now declares its full length, sends 64 KiB and goes silent. The pump stays connected and delivers nothing while the test measures. 12 of 12 authorization-group runs pass, and removing the latch check from detourFetchBlock still fails both latch tests. Co-Authored-By: Claude Opus 5.5 (1M context) --- ...reshableDirectPlayAuthorizationTests.swift | 28 +++++++++++++------ 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index 4c2f68637..33aaf09b0 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -183,7 +183,7 @@ struct RefreshableDirectPlayAuthorizationTests { func detourRefusalLatches() async throws { let provider = ProviderLog() let asksAfterRefusal = Counter() - let server = try Self.origin(total: Self.detourTotal) { _ in true } + let server = try Self.origin(total: Self.detourTotal, stallAt: Self.anchor) { _ in true } defer { server.stop() } let reader = Self.detourReader(server, authorization: HTTPRequestAuthorization { _, rejected in if provider.rotated { @@ -211,7 +211,7 @@ struct RefreshableDirectPlayAuthorizationTests { @Test("a 401 on a backward read is retried once with the refreshed credential") func detourUnauthorizedRetriedWithFreshCredential() async throws { let provider = ProviderLog() - let server = try Self.origin(total: Self.detourTotal) { + let server = try Self.origin(total: Self.detourTotal, stallAt: Self.anchor) { $0.range != Self.detourRange || $0.authorization == "Bearer fresh" } defer { server.stop() } @@ -234,7 +234,7 @@ struct RefreshableDirectPlayAuthorizationTests { @Test("an unchanged credential after a 401 on a backward read fails it, and a repeat read too") func detourUnchangedCredentialFailsTheRead() async throws { let provider = ProviderLog() - let server = try Self.origin(total: Self.detourTotal) { $0.range != Self.detourRange } + let server = try Self.origin(total: Self.detourTotal, stallAt: Self.anchor) { $0.range != Self.detourRange } defer { server.stop() } let reader = Self.detourReader(server, authorization: HTTPRequestAuthorization { _, rejected in provider.answer(rejected: rejected) { _ in "Bearer stale" } @@ -278,14 +278,18 @@ struct RefreshableDirectPlayAuthorizationTests { // MARK: - Support /// A ranged origin over `total` filler bytes. `accepts` decides per request; a refusal is a 401. - private static func origin(total: Int = total, + /// `stallAt`: a range starting there sends 64 KiB of what it declares and then nothing, so the + /// reader keeps a live connection that delivers no more. + private static func origin(total: Int = total, stallAt: Int? = nil, accepts: @escaping @Sendable (ScriptedOriginServer.Recorded) -> Bool) throws -> ScriptedOriginServer { try #require(ScriptedOriginServer { request in guard accepts(request) else { return .init(status: 401, declaredLength: 0) } let (start, end) = Self.bounds(request.range, total: total) - return .init(status: 206, declaredLength: Int64(end - start + 1), - contentRange: "bytes \(start)-\(end)/\(total)", bodyBytes: end - start + 1) + let length = end - start + 1 + return .init(status: 206, declaredLength: Int64(length), + contentRange: "bytes \(start)-\(end)/\(total)", + bodyBytes: start == stallAt ? 64 * 1024 : length) }) } @@ -308,7 +312,9 @@ struct RefreshableDirectPlayAuthorizationTests { } /// The detour tests need a gap behind the window, so their source is large enough to seek more - /// than 8 MiB past a window held to 1 MiB. + /// than 8 MiB past a window held to 1 MiB. A backward read takes the detour only while the pump + /// is connected, and a pump that delivers lifts a latched refusal by design, so the anchored + /// range stalls (`stallAt`): connected, and silent while a test measures. private static let detourTotal = 16 * 1024 * 1024 /// The 4 MiB detour block that holds byte 9 MiB. private static let detourRange = "bytes=8388608-12582911" @@ -320,11 +326,15 @@ struct RefreshableDirectPlayAuthorizationTests { windowHighWater: 1024 * 1024) } - /// Opens, then seeks far enough forward that the pump re-anchors at 12 MiB. + /// Where the detour tests re-anchor the pump: more than 8 MiB past the open's window, so the + /// seek reconnects there instead of reading forward to it. + private static let anchor = 12 * 1024 * 1024 + + /// Opens, then seeks so the pump re-anchors at `anchor`. private static func anchorPastTheHead(_ reader: AVIOReader) throws { try reader.open() _ = read(reader, count: openBytes) - _ = reader.seek(offset: 12 * 1024 * 1024, whence: SEEK_SET) + _ = reader.seek(offset: Int64(anchor), whence: SEEK_SET) _ = read(reader, count: 16 * 1024) } From 39587bc6686566738e39b8d9f861ad63ecb9c7da Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Mon, 5 Oct 2026 14:56:30 -0400 Subject: [PATCH 9/9] test(auth): run the pool-parking resolver test outside the authorization group resolverRunsWhileThePoolIsParked parks twice the cooperative pool's width to prove the resolver no longer needs that pool. It ran in the authorization group, which exists to keep the relay suites' short deadlines away from blocking work. Even a pool parked for milliseconds delays those suites' async work. With the test, the group failed in 3 of 68 runs; without it, in 1 of 40, the same rate as the base branch. The test moves into its own suite, AuthorizationResolverExecutorTests, which runs in the main group. Its only bound is the authorizer's 5 s timeout, which main-group load cannot reach once the resolver is off the pool. It needs no CI step of its own, so ci.yml and CONTRIBUTING.md are unchanged. It is the same test: callers parked on twice the pool width, with a resolver that awaits an actor. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../AuthorizationResolverExecutorTests.swift | 40 +++++++++++++++++++ ...reshableDirectPlayAuthorizationTests.swift | 28 ------------- 2 files changed, 40 insertions(+), 28 deletions(-) create mode 100644 Tests/AetherEngineTests/AuthorizationResolverExecutorTests.swift diff --git a/Tests/AetherEngineTests/AuthorizationResolverExecutorTests.swift b/Tests/AetherEngineTests/AuthorizationResolverExecutorTests.swift new file mode 100644 index 000000000..1cef1787d --- /dev/null +++ b/Tests/AetherEngineTests/AuthorizationResolverExecutorTests.swift @@ -0,0 +1,40 @@ +// Review of PR #12: every open the engine starts runs `Demuxer.open` inside a `Task.detached`, so it +// parks a cooperative-pool thread while it waits for the `HTTPRequestAuthorization` resolver. When +// the resolver needed that same pool, opens that occupied every pool thread left it nowhere to run, +// and each one failed at its bound instead of starting. The resolver now runs on an engine-owned +// executor, so parked callers cannot starve it. +// +// Deliberately NOT in the authorization test group: parking the whole pool, even for milliseconds, +// delays the async work of the deadline-sensitive relay suites that group isolates. Its own bound is +// five seconds, which the main group's blocking work cannot reach once the resolver is off the pool. +import Foundation +import Testing +@testable import AetherEngine + +@Suite("Authorization resolvers run off the cooperative pool", .timeLimit(.minutes(1))) +struct AuthorizationResolverExecutorTests { + /// Twice the pool's width of parked callers, and a resolver that hops onto an actor, is the + /// starved state on any machine: before the fix 22 of 32 callers timed out. + @Test("callers parked on every cooperative thread still get the resolver's answer") + func resolverRunsWhileThePoolIsParked() async throws { + let store = TokenStore() + let authorizer = SourceRequestAuthorizer( + HTTPRequestAuthorization { _, _ in ["Authorization": await store.current()] }, + sourceURL: URL(string: "http://127.0.0.1/movie.mkv")!, timeout: 5) + let callers = ProcessInfo.processInfo.activeProcessorCount * 2 + + let answered = await withTaskGroup(of: Bool.self) { group in + for _ in 0.. String { "Bearer pooled" } +} diff --git a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift index 33aaf09b0..0c589a6aa 100644 --- a/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -252,29 +252,6 @@ struct RefreshableDirectPlayAuthorizationTests { #expect(provider.rejections == ["Bearer stale"]) } - /// Review of PR #12: every open the engine starts runs `Demuxer.open` inside a `Task.detached`, - /// so it parks a cooperative-pool thread while it waits for the resolver. When the resolver - /// needed that same pool, opens that occupied every pool thread left it nowhere to run, and each - /// one failed at its bound instead of starting. Twice the pool's width of parked callers, and a - /// resolver that hops onto an actor, is that state on any machine. - @Test("callers parked on every cooperative thread still get the resolver's answer") - func resolverRunsWhileThePoolIsParked() async throws { - let store = TokenStore() - let authorizer = SourceRequestAuthorizer( - HTTPRequestAuthorization { _, _ in ["Authorization": await store.current()] }, - sourceURL: URL(string: "http://127.0.0.1/movie.mkv")!, timeout: 5) - let callers = ProcessInfo.processInfo.activeProcessorCount * 2 - - let answered = await withTaskGroup(of: Bool.self) { group in - for _ in 0.. String { "Bearer pooled" } -} - /// What the provider was asked. `answer` numbers every call from 1 and records the Authorization of /// each set of rejected headers it was handed. private final class ProviderLog: @unchecked Sendable {