diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index cd6d8193..48d7c25a 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 c91fd1c8..af7fa1be 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,10 +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`. 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. ### 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/CONTRIBUTING.md b/CONTRIBUTING.md index dacf7966..a900dcc1 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 f79aae82..d035904c 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 896f7935..18669c74 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 03aecd5f..9c6950fc 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 522c46ea..d13b342d 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,33 @@ 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 + /// 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) } - 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. `credentialed` is what + /// this request carries. + @discardableResult + private func applyExtraHeaders(_ request: inout URLRequest, + base: [String: String]? = nil) throws -> RedirectHeaderPolicy.Headers { + let headers = try headers(for: request.url, base: base) + for (name, value) in headers.credentialed { request.setValue(value, forHTTPHeaderField: name) } + return headers } func open() throws { @@ -1118,8 +1156,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 +1171,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 +1194,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 +1338,71 @@ 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 { + latchAuthorizationRefusal() + 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 (status 0), or a 401 stood. Latches until a generation + /// delivers again. + private func latchAuthorizationRefusal(status: Int = 0) { + winCond.lock() + authorizationRecovery = .refused(status: status) + 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) + 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; @@ -1340,6 +1447,9 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { close() throw AVIOReaderError.transportSecurityFailed(code: tlsCode) } + if status == 0, let refusedStatus = latchedAuthorizationRefusal() { + throw openFailureForAuthorization(status: refusedStatus) + } guard status != 0 else { return } EngineLog.emit( "[AVIOReader] \(label) source refused: HTTP \(status); failing the open typed", @@ -1463,6 +1573,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 +1618,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 @@ -2037,6 +2149,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) @@ -2103,10 +2218,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 @@ -2137,8 +2263,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) @@ -2199,6 +2325,15 @@ 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): + return failReadForAuthorization(at: frontier, status: refusedStatus, totalRead: totalRead) + 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 +2542,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. @@ -2421,8 +2605,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`. @@ -2458,6 +2648,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) @@ -2474,41 +2666,64 @@ 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) 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 - applyExtraHeaders(&request) - do { - let (data, response) = try syncRequest(request, 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.credentialed) { + 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 } } @@ -2594,7 +2809,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 +2834,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 +2857,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 +3117,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 +3142,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: RedirectHeaderPolicy.Headers + 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 +3164,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 +3175,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 +3193,7 @@ final class AVIOReader: AVIOProvider, @unchecked Sendable { return } activeTransfer = transfer + connSentHeaders = sentHeaders.credentialed winCond.unlock() transfer.startTransfer() @@ -3132,6 +3361,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 +3628,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: RedirectHeaderPolicy.Headers + do { + sentHeaders = try applyExtraHeaders(&request) + } catch { + latchAuthorizationRefusal() + 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 +3661,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 +4005,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 +4022,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 +4058,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 +4096,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: 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. // Scoped to the call: unlike the pump's, this request's life IS this function's. @@ -4059,7 +4310,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 @@ -4182,7 +4433,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 @@ -4195,7 +4446,7 @@ private func redirectPreservingHeaders( request, originalURL: task.originalRequest?.url, originalRange: task.originalRequest?.value(forHTTPHeaderField: "Range"), - extraHeaders: extraHeaders) + headers: extraHeaders) } // MARK: - Persistent Read Delegate @@ -4345,7 +4596,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. @@ -4355,7 +4606,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 @@ -4462,7 +4713,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. @@ -4476,7 +4727,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 } @@ -4584,10 +4835,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, @@ -4660,12 +4911,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 } @@ -4812,7 +5063,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)? @@ -4821,7 +5072,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 } @@ -4941,6 +5192,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 +5206,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 7465a04b..78e87629 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/Demuxer/HeldSourceConnection.swift b/Sources/AetherEngine/Demuxer/HeldSourceConnection.swift index b24fbf48..ba634151 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 18532626..d0e9b465 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/HLSOriginRelay.swift b/Sources/AetherEngine/Network/HLSOriginRelay.swift index 534adabd..33f3d5b1 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 22d5fb1e..6f92e6d0 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. 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 @@ -16,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 @@ -38,14 +45,27 @@ 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 } } -/// 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>? @@ -55,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 { @@ -88,3 +113,75 @@ final class HTTPAuthorizationWait: @unchecked Sendable { running?.cancel() } } + +/// 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. +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 c11d05a5..6f51a08d 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 617085ea..c2016c67 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 0fdbca27..58133f26 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/AuthorizationResolverExecutorTests.swift b/Tests/AetherEngineTests/AuthorizationResolverExecutorTests.swift new file mode 100644 index 00000000..1cef1787 --- /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/Issue255BodyReserveTests.swift b/Tests/AetherEngineTests/Issue255BodyReserveTests.swift index bf140bd7..b0bc22c1 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/Issue377HeldConnectionTests.swift b/Tests/AetherEngineTests/Issue377HeldConnectionTests.swift index 4a9c63db..7f9fddcd 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 41d8a6fb..89943d55 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 new file mode 100644 index 00000000..0c589a6a --- /dev/null +++ b/Tests/AetherEngineTests/RefreshableDirectPlayAuthorizationTests.swift @@ -0,0 +1,383 @@ +// `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 }) + } + + /// 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. + + @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() + 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 { + 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) } + // 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") + } + + @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, stallAt: Self.anchor) { + $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 it, and a repeat read too") + func detourUnchangedCredentialFailsTheRead() async throws { + let provider = ProviderLog() + 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" } + }) + 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) } + 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"]) + } + + // MARK: - Support + + /// A ranged origin over `total` filler bytes. `accepts` decides per request; a refusal is a 401. + /// `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) + let length = end - start + 1 + return .init(status: 206, declaredLength: Int64(length), + contentRange: "bytes \(start)-\(end)/\(total)", + bodyBytes: start == stallAt ? 64 * 1024 : length) + }) + } + + 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) + 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)) + } + + /// 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. 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" + + 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) + } + + /// 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: Int64(anchor), 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) + 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() + } + } +} + +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 } } +} + +/// 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 4d04f8d4..635a8dd8 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 @@ -167,7 +169,29 @@ 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. 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 +`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. 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 @@ -967,7 +991,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. |