refactor: HttpClient

Replaces layering with an inline request pipeline, and transfer queue. This is
meant to simplify the code, reduce footguns, and make future enhancements easier
to implement (e.g. speculative parsing (which requires streaming to fully
leverage)).

Previously, HttpClient implemented deferring as a layer which required special
pumping at various callsites (https://github.com/lightpanda-io/browser/pull/2855,
https://github.com/lightpanda-io/browser/pull/2843, ...). In this new approach,
deferring is built-into the HttpClient/Transfer's flow. Specifically, Transfers
now maintain a queue of events (start, header, data, end, err) which are
dispatched in HttpClient.tick. The result is that JS callbacks are never
executed in the same stack that initiated the I/O, without needing guards or any
external intervention.

tTwo other benefits come from this. The first is that reentrant libcurl is
eliminated. Instead of "libcurl -> callback", it's now "libcurl -> transfer
event queue THEN  tick -> callback" (we don't have to wait until the NEXT tick, we
can just do it later in the tick). HttpClient still has to guard against libcurl
reentrancy, but only because of how WebSocket is implemented, and we should be
able to unify WebSockets to use an event queue too in a follow up PR (which will
eliminate a bunch of guard code).

The transfer queue should also be useful to re-implement streaming, since a
data chunk is just an event in the transfer's event queue. For now, I kept it
as a single buffered event to minimize the change. But since speculative parsing
depends on this, and speculative parsing seems to be the next major performance
tweak we can make, we need to re-introduce streaming.

The other change is the removal of all other layers in favor of a pipeline. This
works well with the existing Transfer.park mechanism, where a parked Transfer
can restart the pipeline for a transfer in an arbitrary point (not as fancy as
it sounds given how simple the flow is). The fallout from this is that we're no
longer creating/wrapping contexts and callbacks: whatever the request was
configured with is all we need.

Because of this, HttpClient.Response is removed. There are no intermediary
responses and no changing context, everything is just the Transfer.

A smaller change is the addition of newRequest + transfer.submit(). The one-shot
HttpClient.request and HttpClient.requestT still exist, but this explicit create
+ submit has some advantage. First, callers can use the transfer.arena (e.g.
Frame using the transfer's arena to set the Referrer header). Second, callers
can holds Transfer immediately, rather than waiting for their startCallback to
be fired. An abort on an XMLHttpRequest called before the start of the transfer
no longer silently fails.
This commit is contained in:
Karl Seguin committed 2026-07-10 07:35:36 +08:00
1 parent 9500c6a653
commit 2eab4d2630
29 files changed
+1772 -2680

No files matched your search

+1 -3
View File
@@ -21,8 +21,7 @@ const lp = @import("lightpanda");
const js = @import("browser/js/js.zig");
const Frame = @import("browser/Frame.zig");
const Transfer = @import("browser/HttpClient.zig").Transfer;
const Response = @import("browser/HttpClient.zig").Response;
const Transfer = @import("network/HttpClient.zig").Transfer;
const log = lp.log;
const Execution = js.Execution;
@@ -224,7 +223,6 @@ pub const ResponseData = struct {
pub const ResponseHeaderDone = struct {
transfer: *Transfer,
response: *const Response,
};
pub const RequestDone = struct {
+1 -1
View File
@@ -28,7 +28,7 @@ const Watchdog = @import("../Watchdog.zig");
const Session = @import("Session.zig");
const Selector = @import("webapi/selector/Selector.zig");
const Viewport = @import("Viewport.zig");
const HttpClient = @import("HttpClient.zig");
const HttpClient = @import("../network/HttpClient.zig");
const PermissionState = @import("webapi/Permissions.zig").State;
const ArenaPool = App.ArenaPool;
+65 -64
View File
@@ -60,7 +60,7 @@ const popover = @import("webapi/element/popover.zig");
const slotting = @import("webapi/element/slotting.zig");
const NavigationKind = @import("webapi/navigation/root.zig").NavigationKind;
const HttpClient = @import("HttpClient.zig");
const HttpClient = @import("../network/HttpClient.zig");
const sys_url = @import("../sys/url.zig");
const timestamp = @import("../datetime.zig").timestamp;
@@ -597,16 +597,16 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
self._load_state = .parsing;
self._last_navigate_error = null;
const req_id = self._session.browser.http_client.nextReqId();
log.info(.frame, "navigate", .{
.url = request_url,
.method = opts.method,
.reason = opts.reason,
.body = opts.body != null,
.req_id = req_id,
.type = self._type,
});
const http_client = &session.browser.http_client;
// Handle synthetic navigations: about:blank and blob: URLs
const is_about_blank = std.mem.eql(u8, "about:blank", request_url);
const is_blob = !is_about_blank and std.mem.startsWith(u8, request_url, "blob:");
@@ -670,6 +670,10 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
};
}
// No real request is made, but CDP still correlates the navigation
// events by request id, so consume one.
const req_id = http_client.incrReqId();
session.notification.dispatch(.frame_navigate, &.{
.opts = opts,
.req_id = req_id,
@@ -694,9 +698,6 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
.timestamp = timestamp(.monotonic),
});
// force next request id manually b/c we won't create a real req.
_ = session.browser.http_client.incrReqId();
if (self.parent == null) {
session.navigation._current_navigation_kind = opts.kind;
try session.navigation.commitNavigation(self);
@@ -706,8 +707,6 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
return;
}
const http_client = &session.browser.http_client;
self._http_status = null;
self._http_headers = .empty;
@@ -719,7 +718,6 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
};
self.origin = try URL.getOrigin(self.arena, self.url);
self._req_id = req_id;
self._navigated_options = .{
.cdp_id = opts.cdp_id,
.reason = opts.reason,
@@ -728,14 +726,37 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
.header = if (opts.header) |h| try self.arena.dupeZ(u8, h) else null,
};
var headers = try http_client.newHeaders();
try headers.add(lp.Config.HttpHeaders.navigation_accept);
if (opts.header) |hdr| {
try headers.add(hdr);
}
if (opts.referer) |ref| {
const ref_header = try std.mem.concatWithSentinel(self.arena, u8, &.{ "Referer: ", ref }, 0);
try headers.add(ref_header);
const transfer = try http_client.newRequest(.{
.ctx = self,
.url = self.url,
.frame_id = self._frame_id,
.loader_id = self._loader_id,
.method = opts.method,
.body = opts.body,
.cookie_jar = &session.cookie_jar,
.cookie_origin = opts.initiator_url orelse self.url,
.resource_type = .document,
.notification = self._session.notification,
.header_callback = frameHeaderDoneCallback,
.data_callback = frameDataCallback,
.done_callback = frameDoneCallback,
.error_callback = frameErrorCallback,
// The frame tracks its navigation by id (_req_id), never by pointer.
.shutdown_callback = HttpClient.noopShutdown,
}, &self._http_owner);
self._req_id = transfer.id;
{
// Ours until submit; clean up if header setup fails.
errdefer transfer.deinit();
try transfer.req.headers.add(lp.Config.HttpHeaders.navigation_accept);
if (opts.header) |hdr| {
try transfer.req.headers.add(hdr);
}
if (opts.referer) |ref| {
const ref_header = try std.mem.concatWithSentinel(transfer.arena, u8, &.{ "Referer: ", ref }, 0);
try transfer.req.headers.add(ref_header);
}
}
// A root navigation issued against a pending Page (i.e. one allocated by
@@ -751,7 +772,7 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
session.notification.dispatch(.frame_navigate, &.{
.opts = opts,
.url = self.url,
.req_id = req_id,
.req_id = transfer.id,
.frame_id = self._frame_id,
.loader_id = self._loader_id,
.timestamp = timestamp(.monotonic),
@@ -763,23 +784,7 @@ pub fn navigate(self: *Frame, request_url: [:0]const u8, opts: NavigateOpts) !vo
session.navigation._current_navigation_kind = opts.kind;
self.makeRequest(.{
.ctx = self,
.url = self.url,
.frame_id = self._frame_id,
.loader_id = self._loader_id,
.method = opts.method,
.headers = headers,
.body = opts.body,
.cookie_jar = &session.cookie_jar,
.cookie_origin = opts.initiator_url orelse self.url,
.resource_type = .document,
.notification = self._session.notification,
.header_callback = frameHeaderDoneCallback,
.data_callback = frameDataCallback,
.done_callback = frameDoneCallback,
.error_callback = frameErrorCallback,
}) catch |err| {
transfer.submit() catch |err| {
log.err(.frame, "navigate request", .{ .url = self.url, .err = err, .type = self._type });
return err;
};
@@ -962,6 +967,11 @@ pub fn makeRequest(self: *Frame, req: HttpClient.Request) !void {
return self._session.browser.http_client.request(req, &self._http_owner);
}
// Two-phase variant; see HttpClient.newRequest for the ownership contract.
pub fn newRequest(self: *Frame, req: HttpClient.Request) !*HttpClient.Transfer {
return self._session.browser.http_client.newRequest(req, &self._http_owner);
}
// Synchronously abort every transfer and WebSocket owned by this frame
// and all of its descendants.
pub fn abortTransfers(self: *Frame) void {
@@ -970,8 +980,6 @@ pub fn abortTransfers(self: *Frame) void {
}
const http_client = &self._session.browser.http_client;
http_client.abortOwner(&self._http_owner);
// abortOwner misses deferred contexts whose transfer already completed.
http_client.deferring_layer.cancelFrame(self._frame_id);
}
pub fn documentIsLoaded(self: *Frame) void {
@@ -1144,8 +1152,8 @@ fn notifyParentLoadComplete(self: *Frame) void {
parent.iframeCompletedLoading(self.iframe.?, self._delays_parent_load);
}
fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
var self: *Frame = @ptrCast(@alignCast(response.ctx));
fn frameHeaderDoneCallback(transfer: *HttpClient.Transfer) !HttpClient.Transfer.HeaderResult {
var self: *Frame = @ptrCast(@alignCast(transfer.req.ctx));
// Commit point for a pending root navigation. The session has been
// holding the OLD page alive during the round-trip; now that response
@@ -1157,7 +1165,7 @@ fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResu
try self._session.commitPendingPage(self._page);
}
const response_url = response.url();
const response_url = transfer.req.url;
if (std.mem.eql(u8, response_url, self.url) == false) {
// would be different than self.url in the case of a redirect
self.url = try self.arena.dupeZ(u8, response_url);
@@ -1169,7 +1177,7 @@ fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResu
// Page.reload doesn't re-POST form data to the redirect target. Conservative
// default — 307/308 technically preserve the method per RFC 7231, but
// resubmitting form data is the more dangerous failure mode.
if ((response.redirectCount() orelse 0) > 0) {
if ((transfer.redirectCount() orelse 0) > 0) {
if (self._navigated_options) |*no| {
no.method = .GET;
no.body = null;
@@ -1186,14 +1194,14 @@ fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResu
if (comptime IS_DEBUG) {
log.debug(.frame, "navigate header", .{
.url = self.url,
.status = response.status(),
.content_type = response.contentType(),
.status = transfer.responseStatus(),
.content_type = transfer.contentType(),
.type = self._type,
});
}
self._http_status = response.status();
var it = response.headerIterator();
self._http_status = transfer.responseStatus();
var it = transfer.responseHeaderIterator();
while (it.next()) |hdr| {
try self._http_headers.append(self.arena, .{
.name = try self.arena.dupe(u8, hdr.name),
@@ -1218,7 +1226,7 @@ fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResu
// If the response is a file download, stream its body to disk instead of
// parsing it as a page. This sets _parse_state to .download, which the
// data/done callbacks below special-case.
_ = try self.maybeStartDownload(response);
_ = try self.maybeStartDownload(transfer);
return .proceed;
}
@@ -1227,7 +1235,7 @@ fn frameHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResu
// treated as a download when Browser.setDownloadBehavior opted in
// (allow/allowAndName) and the response carries Content-Disposition: attachment.
// See issue #2701.
fn maybeStartDownload(self: *Frame, response: HttpClient.Response) !bool {
fn maybeStartDownload(self: *Frame, transfer: *HttpClient.Transfer) !bool {
const session = self._session;
switch (session.download_behavior) {
.allow, .allow_and_name => {},
@@ -1235,7 +1243,7 @@ fn maybeStartDownload(self: *Frame, response: HttpClient.Response) !bool {
}
const disposition: HttpClient.Header = blk: {
var it = response.headerIterator();
var it = transfer.responseHeaderIterator();
while (it.next()) |hdr| {
if (std.ascii.eqlIgnoreCase(hdr.name, "content-disposition")) {
break :blk hdr;
@@ -1281,7 +1289,7 @@ fn maybeStartDownload(self: *Frame, response: HttpClient.Response) !bool {
return false;
};
const total: ?u64 = if (response.contentLength()) |cl| cl else null;
const total: ?u64 = if (transfer.getContentLength()) |cl| cl else null;
self._parse_state = .{ .download = .{
.guid = guid,
@@ -1364,14 +1372,14 @@ fn isUtf16Encoding(charset: []const u8) bool {
return std.mem.eql(u8, charset, "UTF-16LE") or std.mem.eql(u8, charset, "UTF-16BE");
}
fn frameDataCallback(response: HttpClient.Response, data: []const u8) !void {
var self: *Frame = @ptrCast(@alignCast(response.ctx));
fn frameDataCallback(transfer: *HttpClient.Transfer, data: []const u8) !void {
var self: *Frame = @ptrCast(@alignCast(transfer.req.ctx));
if (self._parse_state == .pre) {
// we lazily do this, because we might need the first chunk of data
// to sniff the content type
var mime: Mime = blk: {
if (response.contentType()) |ct| {
if (transfer.contentType()) |ct| {
break :blk try Mime.parse(ct);
}
break :blk Mime.sniff(data);
@@ -2096,18 +2104,10 @@ pub fn loadExternalStylesheet(self: *Frame, link: *Element.Html.Link, href: []co
const http_client = &session.browser.http_client;
// `syncRequest` below registers a blocking request for this frame, which
// makes the DeferringLayer hold back the completion callbacks of every
// OTHER in-flight transfer for the frame (e.g. a `<script defer>` still
// loading) so they don't run JS while we're on the parser stack. Those
// deferred completions must be flushed once the sync fetch returns —
// otherwise a `<script defer>` whose fetch finishes during this window is
// left at `complete == false` forever, the deferred-script queue never
// drains, and `documentIsLoaded` (readyState -> "interactive",
// DOMContentLoaded, the load event) never fires. The blocking-`<script>`
// path (ScriptManager.addFromElement) and the worker path already flush;
// the external-stylesheet path must too.
defer http_client.deferring_layer.flushFrame(self._frame_id);
// `syncRequest` below registers a blocking request for this frame; the
// client's dispatcher holds back every OTHER transfer's callbacks for
// the frame while it's registered (they'd run JS on the parser's stack)
// and delivers them on the next tick after the sync fetch returns.
var headers = try http_client.newHeaders();
try headers.add("Accept: text/css,*/*;q=0.1");
@@ -2135,6 +2135,7 @@ pub fn loadExternalStylesheet(self: *Frame, link: *Element.Html.Link, href: []co
.cookie_origin = self.url,
.resource_type = .stylesheet,
.notification = session.notification,
.shutdown_callback = HttpClient.noopShutdown, // syncRequest installs its own
}) catch |err| {
log.warn(.http, "external stylesheet fetch", .{ .err = err, .url = resolved });
return self.fireElementEvent(element, comptime .wrap("error"));
+4 -4
View File
@@ -22,7 +22,7 @@ const lp = @import("lightpanda");
const js = @import("js/js.zig");
const Browser = @import("Browser.zig");
const Session = @import("Session.zig");
const HttpClient = @import("HttpClient.zig");
const HttpClient = @import("../network/HttpClient.zig");
const Node = @import("webapi/Node.zig");
const Selector = @import("webapi/selector/Selector.zig");
@@ -221,12 +221,12 @@ fn _tick(self: *Runner, comptime is_cdp: bool, timeout_ms: u32, conditions: []Wa
}
const http_active = http_client.http_active;
const http_next_tick = http_client.next_tick_count;
const total_http_activity = http_active + http_next_tick + http_client.interception_layer.intercepted;
const http_buffered = http_client.dispatch_count;
const total_http_activity = http_active + http_buffered + http_client.intercepted;
const total_network_activity = total_http_activity + http_client.ws_active;
const ms_to_next_macrotask = browser.msToNextMacrotask();
const network_idle = total_network_activity == 0 and http_client.queue.first == null and http_client.ready_queue.first == null;
const network_idle = total_network_activity == 0 and http_client.pending_queue.first == null and http_client.ready_queue.first == null;
const is_done = ms_to_next_macrotask == null and network_idle;
// _we_ have nothing to run, but v8 is working on background tasks. We'll
+5 -4
View File
@@ -20,7 +20,7 @@ const std = @import("std");
const lp = @import("lightpanda");
const builtin = @import("builtin");
const HttpClient = @import("HttpClient.zig");
const HttpClient = @import("../network/HttpClient.zig");
const js = @import("js/js.zig");
const URL = @import("URL.zig");
@@ -358,6 +358,7 @@ pub fn addFromElement(self: *ScriptManager, comptime from_parser: bool, script_e
.cookie_origin = frame.url,
.resource_type = .script,
.notification = frame._session.notification,
.shutdown_callback = HttpClient.noopShutdown, // syncRequest installs its own
});
script.source = .{ .remote = response.body };
@@ -386,6 +387,9 @@ pub fn addFromElement(self: *ScriptManager, comptime from_parser: bool, script_e
.data_callback = Script.dataCallback,
.done_callback = Script.doneCallback,
.error_callback = Script.errorCallback,
// Nothing holds the transfer; teardown cleanup runs through
// the manager's script lists.
.shutdown_callback = HttpClient.noopShutdown,
});
}
@@ -396,9 +400,6 @@ pub fn addFromElement(self: *ScriptManager, comptime from_parser: bool, script_e
return;
}
// This will flush any deferred scripts.
defer self.base.client.deferring_layer.flushFrame(self.base.owner.frameId());
if (script.status < 200 or script.status > 299) {
log.info(.http, "script load error", .{ .status = script.status });
script.executeCallback(comptime .wrap("error"));
+49 -51
View File
@@ -20,8 +20,8 @@ const std = @import("std");
const lp = @import("lightpanda");
const builtin = @import("builtin");
const HttpClient = @import("HttpClient.zig");
const http = @import("../network/http.zig");
const HttpClient = @import("../network/HttpClient.zig");
const js = @import("js/js.zig");
const Session = @import("Session.zig");
@@ -454,6 +454,7 @@ pub fn getAsyncImport(self: *ScriptManagerBase, url: [:0]const u8, cb: ImportAsy
.data_callback = Script.dataCallback,
.done_callback = Script.doneCallback,
.error_callback = Script.errorCallback,
.shutdown_callback = Script.shutdownCallback,
}) catch |err| {
self.async_scripts.remove(&script.node);
return err;
@@ -660,19 +661,19 @@ pub const Script = struct {
self.manager.releaseArena(self.arena);
}
pub fn startCallback(response: HttpClient.Response) !void {
log.debug(.http, "script fetch start", .{ .req = response });
pub fn startCallback(transfer: *HttpClient.Transfer) !void {
log.debug(.http, "script fetch start", .{ .req = transfer });
}
pub fn headerCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
const self: *Script = @ptrCast(@alignCast(response.ctx));
pub fn headerCallback(transfer: *HttpClient.Transfer) !HttpClient.Transfer.HeaderResult {
const self: *Script = @ptrCast(@alignCast(transfer.req.ctx));
self.status = response.status().?;
if (response.status() != 200) {
self.status = transfer.responseStatus().?;
if (transfer.responseStatus() != 200) {
log.info(.http, "script header", .{
.req = response,
.status = response.status(),
.content_type = response.contentType(),
.req = transfer,
.status = transfer.responseStatus(),
.content_type = transfer.contentType(),
});
return .abort;
@@ -680,64 +681,61 @@ pub const Script = struct {
if (comptime IS_DEBUG) {
log.debug(.http, "script header", .{
.req = response,
.status = response.status(),
.content_type = response.contentType(),
.req = transfer,
.status = transfer.responseStatus(),
.content_type = transfer.contentType(),
});
}
switch (response.inner) {
.transfer => |transfer| {
// temp debug, trying to figure out why the next assert sometimes
// fails. Is the buffer just corrupt or is headerCallback really
// being called twice?
lp.assert(self.header_callback_called == false, "ScriptManagerBase.Header recall", .{
.m = @tagName(std.meta.activeTag(self.extra)),
.a1 = self.debug_transfer_id,
.a2 = self.debug_transfer_tries,
.a3 = self.debug_transfer_aborted,
.a4 = self.debug_transfer_bytes_received,
.a5 = self.debug_transfer_notified_fail,
.a8 = self.debug_transfer_auth_challenge,
.a9 = self.debug_transfer_easy_id,
.b1 = transfer.id,
.b2 = transfer._tries,
.b3 = transfer.state == .aborted,
.b4 = transfer.res.bytes_received,
.b5 = transfer._notified_fail,
.b8 = transfer._auth_challenge != null,
.b9 = if (transfer._conn) |c| @intFromPtr(c._easy) else 0,
});
self.header_callback_called = true;
self.debug_transfer_id = transfer.id;
self.debug_transfer_tries = transfer._tries;
self.debug_transfer_aborted = transfer.state == .aborted;
self.debug_transfer_bytes_received = transfer.res.bytes_received;
self.debug_transfer_notified_fail = transfer._notified_fail;
self.debug_transfer_auth_challenge = transfer._auth_challenge != null;
self.debug_transfer_easy_id = if (transfer._conn) |c| @intFromPtr(c._easy) else 0;
},
else => {},
{
// temp debug, trying to figure out why the next assert sometimes
// fails. Is the buffer just corrupt or is headerCallback really
// being called twice?
lp.assert(self.header_callback_called == false, "ScriptManagerBase.Header recall", .{
.m = @tagName(std.meta.activeTag(self.extra)),
.a1 = self.debug_transfer_id,
.a2 = self.debug_transfer_tries,
.a3 = self.debug_transfer_aborted,
.a4 = self.debug_transfer_bytes_received,
.a5 = self.debug_transfer_notified_fail,
.a8 = self.debug_transfer_auth_challenge,
.a9 = self.debug_transfer_easy_id,
.b1 = transfer.id,
.b2 = transfer._tries,
.b3 = transfer.state == .aborted,
.b4 = transfer.res.bytes_received,
.b5 = transfer._notified_fail,
.b8 = transfer._auth_challenge != null,
.b9 = if (transfer._conn) |c| @intFromPtr(c._easy) else 0,
});
self.header_callback_called = true;
self.debug_transfer_id = transfer.id;
self.debug_transfer_tries = transfer._tries;
self.debug_transfer_aborted = transfer.state == .aborted;
self.debug_transfer_bytes_received = transfer.res.bytes_received;
self.debug_transfer_notified_fail = transfer._notified_fail;
self.debug_transfer_auth_challenge = transfer._auth_challenge != null;
self.debug_transfer_easy_id = if (transfer._conn) |c| @intFromPtr(c._easy) else 0;
}
lp.assert(self.source.remote.capacity == 0, "ScriptManagerBase.Header buffer", .{ .capacity = self.source.remote.capacity });
var buffer: std.ArrayList(u8) = .empty;
if (response.contentLength()) |cl| {
if (transfer.getContentLength()) |cl| {
try buffer.ensureTotalCapacity(self.arena, cl);
}
self.source = .{ .remote = buffer };
return .proceed;
}
pub fn dataCallback(response: HttpClient.Response, data: []const u8) !void {
const self: *Script = @ptrCast(@alignCast(response.ctx));
self._dataCallback(response, data) catch |err| {
log.err(.http, "SM.dataCallback", .{ .err = err, .transfer = response, .len = data.len });
pub fn dataCallback(transfer: *HttpClient.Transfer, data: []const u8) !void {
const self: *Script = @ptrCast(@alignCast(transfer.req.ctx));
self._dataCallback(transfer, data) catch |err| {
log.err(.http, "SM.dataCallback", .{ .err = err, .transfer = transfer, .len = data.len });
return err;
};
}
fn _dataCallback(self: *Script, _: HttpClient.Response, data: []const u8) !void {
fn _dataCallback(self: *Script, _: *HttpClient.Transfer, data: []const u8) !void {
try self.source.remote.appendSlice(self.arena, data);
}
+8 -1
View File
@@ -33,7 +33,7 @@ const Scheduler = @import("Scheduler.zig");
const Page = @import("../Page.zig");
const Session = @import("../Session.zig");
const Factory = @import("../Factory.zig");
const HttpClient = @import("../HttpClient.zig");
const HttpClient = @import("../../network/HttpClient.zig");
const EventManagerBase = @import("../EventManagerBase.zig");
const Event = @import("../webapi/Event.zig");
@@ -102,6 +102,13 @@ pub fn makeRequest(self: *const Execution, req: HttpClient.Request) !void {
};
}
// Two-phase variant; see HttpClient.newRequest for the ownership contract.
pub fn newRequest(self: *const Execution, req: HttpClient.Request) !*HttpClient.Transfer {
return switch (self.js.global) {
inline else => |g| g.newRequest(req),
};
}
pub fn getBroadcastChannels(self: *const Execution) *std.DoublyLinkedList {
return switch (self.js.global) {
inline else => |g| &g._broadcast_channels,
+33 -18
View File
@@ -23,7 +23,7 @@ const js = @import("../js/js.zig");
const URL = @import("../URL.zig");
const Frame = @import("../Frame.zig");
const HttpClient = @import("../HttpClient.zig");
const Transfer = @import("../../network/HttpClient.zig").Transfer;
const EventTarget = @import("EventTarget.zig");
const MessageEvent = @import("event/MessageEvent.zig");
@@ -56,7 +56,7 @@ _url: [:0]const u8,
_type: WorkerType = .classic,
_script_loaded: bool = false,
_script_buffer: std.ArrayList(u8) = .empty,
_http_response: ?HttpClient.Response = null,
_http_transfer: ?*Transfer = null,
// Event handlers
_on_error: ?js.Function.Global = null,
@@ -101,11 +101,9 @@ pub fn init(url: []const u8, options: ?WorkerOptions, frame: *Frame) !*Worker {
return self;
}
const headers = try session.browser.http_client.newHeaders();
frame.makeRequest(.{
const transfer = frame.newRequest(.{
.ctx = self,
.method = .GET,
.headers = headers,
.url = resolved_url,
.frame_id = self._frame_id,
.loader_id = self._loader_id,
@@ -117,11 +115,24 @@ pub fn init(url: []const u8, options: ?WorkerOptions, frame: *Frame) !*Worker {
.data_callback = httpDataCallback,
.done_callback = httpDoneCallback,
.error_callback = httpErrorCallback,
.shutdown_callback = httpShutdownCallback,
}) catch |err| {
log.err(.browser, "Worker request", .{ .url = resolved_url, .err = err });
frame.removeWorker(self);
return err;
};
// Held for deinit's abort; the done, error and shutdown callbacks clear
// it. The shutdown one matters: Frame.deinit's abortOwner kills the
// transfer before it deinits this worker, and deinit must not abort a
// freed transfer.
self._http_transfer = transfer;
transfer.submit() catch |err| {
log.err(.browser, "Worker request", .{ .url = resolved_url, .err = err });
frame.removeWorker(self);
return err;
};
return self;
}
@@ -129,9 +140,9 @@ pub fn init(url: []const u8, options: ?WorkerOptions, frame: *Frame) !*Worker {
// remove from the frame's worker list.
pub fn deinit(self: *Worker) void {
// No pending frame for workers, so we can abort all frames.
if (self._http_response) |res| {
if (self._http_transfer) |res| {
res.abort(error.Abort);
self._http_response = null;
self._http_transfer = null;
}
self._worker_scope.deinit();
self._frame._session.releaseArena(self._arena);
@@ -141,10 +152,10 @@ pub fn asEventTarget(self: *Worker) *EventTarget {
return self._proto;
}
fn httpHeaderCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
const self: *Worker = @ptrCast(@alignCast(response.ctx));
fn httpHeaderCallback(transfer: *Transfer) !Transfer.HeaderResult {
const self: *Worker = @ptrCast(@alignCast(transfer.req.ctx));
const status = response.status() orelse return .abort;
const status = transfer.responseStatus() orelse return .abort;
if (status < 200 or status >= 300) {
log.warn(.browser, "Worker status", .{
.url = self._url,
@@ -153,22 +164,21 @@ fn httpHeaderCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
return .abort;
}
self._http_response = response;
if (response.contentLength()) |cl| {
if (transfer.getContentLength()) |cl| {
try self._script_buffer.ensureTotalCapacity(self._arena, cl);
}
return .proceed;
}
fn httpDataCallback(response: HttpClient.Response, data: []const u8) !void {
const self: *Worker = @ptrCast(@alignCast(response.ctx));
fn httpDataCallback(transfer: *Transfer, data: []const u8) !void {
const self: *Worker = @ptrCast(@alignCast(transfer.req.ctx));
try self._script_buffer.appendSlice(self._arena, data);
}
fn httpDoneCallback(ctx: *anyopaque) !void {
const self: *Worker = @ptrCast(@alignCast(ctx));
self._http_response = null;
self._http_transfer = null;
const url = self._url;
const script = self._script_buffer.items;
@@ -245,9 +255,14 @@ fn loadInitialScript(self: *Worker, script: []const u8) !void {
ls.local.runMacrotasks();
}
fn httpShutdownCallback(ctx: *anyopaque) void {
const self: *Worker = @ptrCast(@alignCast(ctx));
self._http_transfer = null;
}
fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void {
const self: *Worker = @ptrCast(@alignCast(ctx));
self._http_response = null;
self._http_transfer = null;
log.err(.browser, "worker fetch error", .{
.url = self._url,
@@ -297,9 +312,9 @@ fn _fireErrorEvent(self: *Worker, message: []const u8, error_value: ?js.Value.Te
pub fn terminate(self: *Worker) void {
// Abort any pending script fetch
if (self._http_response) |resp| {
if (self._http_transfer) |resp| {
resp.abort(error.Abort);
self._http_response = null;
self._http_transfer = null;
}
}
+7 -3
View File
@@ -29,7 +29,7 @@ const Page = @import("../Page.zig");
const Frame = @import("../Frame.zig");
const Factory = @import("../Factory.zig");
const Session = @import("../Session.zig");
const HttpClient = @import("../HttpClient.zig");
const HttpClient = @import("../../network/HttpClient.zig");
const EventManagerBase = @import("../EventManagerBase.zig");
const ScriptManagerBase = @import("../ScriptManagerBase.zig");
@@ -250,6 +250,11 @@ pub fn makeRequest(self: *WorkerGlobalScope, req: HttpClient.Request) !void {
return self._session.browser.http_client.request(req, &self._http_owner);
}
// Two-phase variant; see HttpClient.newRequest for the ownership contract.
pub fn newRequest(self: *WorkerGlobalScope, req: HttpClient.Request) !*HttpClient.Transfer {
return self._session.browser.http_client.newRequest(req, &self._http_owner);
}
pub fn getSelf(self: *WorkerGlobalScope) *WorkerGlobalScope {
return self;
}
@@ -390,6 +395,7 @@ fn importScript(self: *WorkerGlobalScope, arena: Allocator, url: [:0]const u8) !
.cookie_origin = self.url,
.resource_type = .script,
.notification = session.notification,
.shutdown_callback = HttpClient.noopShutdown, // syncRequest installs its own
}) catch |err| {
log.warn(.http, "importScript", .{ .url = resolved_url, .err = err });
return error.NetworkError;
@@ -400,8 +406,6 @@ fn importScript(self: *WorkerGlobalScope, arena: Allocator, url: [:0]const u8) !
return error.NetworkError;
}
defer http_client.deferring_layer.flushFrame(self._frame_id);
var ls: JS.Local.Scope = undefined;
self.js.localScope(&ls);
defer ls.deinit();
+3 -4
View File
@@ -290,10 +290,9 @@ test "WebApi: HTML.Link external stylesheet" {
}
// Regression: a synchronous external-stylesheet fetch must not strand the
// completion of an in-flight <script defer> (deferred by the blocking-request
// window). Otherwise the deferred-script queue never drains and the document
// is stuck at readyState "loading". See Frame.loadExternalStylesheet's
// flushFrame call.
// completion of an in-flight <script defer> (held back by the blocking-
// request gate during the sync window). Otherwise the deferred-script queue
// never drains and the document is stuck at readyState "loading".
test "WebApi: HTML.Link deferred script then external stylesheet" {
const filter: testing.LogFilter = .init(&.{.http});
defer filter.deinit();
+30 -31
View File
@@ -18,10 +18,10 @@
const std = @import("std");
const lp = @import("lightpanda");
const HttpClient = @import("../../HttpClient.zig");
const js = @import("../../js/js.zig");
const URL = @import("../../URL.zig");
const Transfer = @import("../../../network/HttpClient.zig").Transfer;
const Request = @import("Request.zig");
const Response = @import("Response.zig");
@@ -96,12 +96,7 @@ pub fn init(input: Input, options: ?InitOpts, exec: *const Execution) !js.Promis
.@"same-origin" => if (exec.isSameOrigin(request._url)) &session.cookie_jar else null,
};
// Synchronous failures from request layers (e.g. RobotsLayer returning
// RobotsBlocked when robots.txt is already cached) are dispatched to
// httpErrorCallback by Client.request, which rejects the promise and
// releases response._arena. Propagating the error from here would also
// fire the `errdefer response.deinit` above and double-free the arena.
exec.makeRequest(.{
const transfer = exec.newRequest(.{
.ctx = fetch,
.url = request._url,
.method = request._method,
@@ -118,26 +113,30 @@ pub fn init(input: Input, options: ?InitOpts, exec: *const Execution) !js.Promis
.@"error" => .@"error",
},
.notification = session.notification,
.start_callback = httpStartCallback,
.header_callback = httpHeaderDoneCallback,
.data_callback = httpDataCallback,
.done_callback = httpDoneCallback,
.error_callback = httpErrorCallback,
.shutdown_callback = httpShutdownCallback,
}) catch {};
}) catch {
// OOM-class; nothing was committed and no callback fired.
return resolver.promise();
};
// Held for Response.deinit's abort; the error, shutdown and done
// callbacks clear it.
response._http_transfer = transfer;
// Failures inside submit are dispatched to httpErrorCallback, which
// rejects the promise and releases response._arena. Propagating the
// error from here would also fire the `errdefer response.deinit` above
// and double-free the arena.
transfer.submit() catch {};
return resolver.promise();
}
fn httpStartCallback(response: HttpClient.Response) !void {
const self: *Fetch = @ptrCast(@alignCast(response.ctx));
if (comptime IS_DEBUG) {
log.debug(.http, "request start", .{ .url = self._url, .source = "fetch" });
}
self._response._http_response = response;
}
fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
const self: *Fetch = @ptrCast(@alignCast(response.ctx));
fn httpHeaderDoneCallback(transfer: *Transfer) !Transfer.HeaderResult {
const self: *Fetch = @ptrCast(@alignCast(transfer.req.ctx));
if (self._signal) |signal| {
if (signal._aborted) {
@@ -146,7 +145,7 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
}
const arena = self._response._arena;
if (response.contentLength()) |cl| {
if (transfer.getContentLength()) |cl| {
try self._buf.ensureTotalCapacity(arena, cl);
}
@@ -156,14 +155,14 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
log.debug(.http, "request header", .{
.source = "fetch",
.url = self._url,
.status = response.status(),
.status = transfer.responseStatus(),
});
}
res._status = response.status().?;
res._status_text = std.http.Status.phrase(@enumFromInt(response.status().?)) orelse "";
res._url = try arena.dupeZ(u8, response.url());
res._is_redirected = response.redirectCount().? > 0;
res._status = transfer.responseStatus().?;
res._status_text = std.http.Status.phrase(@enumFromInt(transfer.responseStatus().?)) orelse "";
res._url = try arena.dupeZ(u8, transfer.req.url);
res._is_redirected = transfer.redirectCount().? > 0;
// redirect: "manual" surfaces the unfollowed 3xx as an opaque-redirect
// filtered response: status 0, no headers, no body.
@@ -195,7 +194,7 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
res._type = .basic;
}
var it = response.headerIterator();
var it = transfer.responseHeaderIterator();
while (it.next()) |hdr| {
try res._headers.append(hdr.name, hdr.value, exec);
}
@@ -203,8 +202,8 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
return .proceed;
}
fn httpDataCallback(response: HttpClient.Response, data: []const u8) !void {
const self: *Fetch = @ptrCast(@alignCast(response.ctx));
fn httpDataCallback(transfer: *Transfer, data: []const u8) !void {
const self: *Fetch = @ptrCast(@alignCast(transfer.req.ctx));
// Check if aborted
if (self._signal) |signal| {
@@ -219,7 +218,7 @@ fn httpDataCallback(response: HttpClient.Response, data: []const u8) !void {
fn httpDoneCallback(ctx: *anyopaque) !void {
const self: *Fetch = @ptrCast(@alignCast(ctx));
var response = self._response;
response._http_response = null;
response._http_transfer = null;
response._body = .{ .bytes = self._buf.items };
log.info(.http, "request complete", .{
@@ -249,7 +248,7 @@ fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void {
});
var response = self._response;
response._http_response = null;
response._http_transfer = null;
// Capture this before we reject. Rejection could trigger httpShutdownCallback
// (via a microtask callback). But if we're here, then we'll take care of
@@ -277,7 +276,7 @@ fn httpShutdownCallback(ctx: *anyopaque) void {
if (self._owns_response) {
var response = self._response;
response._http_response = null;
response._http_transfer = null;
response.deinit(self._exec.page);
// Do not access `self` after this point: the Fetch struct was
// allocated from response._arena which has been released.
+5 -5
View File
@@ -22,7 +22,7 @@ const lp = @import("lightpanda");
const js = @import("../../js/js.zig");
const URL = @import("../../URL.zig");
const Page = @import("../../Page.zig");
const HttpClient = @import("../../HttpClient.zig");
const Transfer = @import("../../../network/HttpClient.zig").Transfer;
const Blob = @import("../Blob.zig");
const ReadableStream = @import("../streams/ReadableStream.zig");
@@ -53,7 +53,7 @@ _type: Type,
_status_text: []const u8,
_url: [:0]const u8,
_is_redirected: bool,
_http_response: ?HttpClient.Response = null,
_http_transfer: ?*Transfer = null,
_body_used: bool = false,
const Body = union(enum) {
@@ -197,9 +197,9 @@ pub fn createJson(data: js.Value, opts_: ?InitOpts, exec: *const Execution) !*Re
}
pub fn deinit(self: *Response, page: *Page) void {
if (self._http_response) |resp| {
if (self._http_transfer) |resp| {
resp.abort(error.Abort);
self._http_response = null;
self._http_transfer = null;
}
page.releaseArena(self._arena);
}
@@ -484,7 +484,7 @@ pub fn clone(self: *const Response, exec: *const Execution) !*Response {
._type = self._type,
._is_redirected = self._is_redirected,
._headers = try Headers.init(.{ .obj = self._headers }, exec),
._http_response = null,
._http_transfer = null,
};
return cloned;
}
+1 -1
View File
@@ -26,7 +26,7 @@ const Blob = @import("../Blob.zig");
const URL = @import("../../URL.zig");
const Page = @import("../../Page.zig");
const HttpClient = @import("../../HttpClient.zig");
const HttpClient = @import("../../../network/HttpClient.zig");
const Event = @import("../Event.zig");
const EventTarget = @import("../EventTarget.zig");
+42 -34
View File
@@ -20,8 +20,8 @@ const std = @import("std");
const lp = @import("lightpanda");
const js = @import("../../js/js.zig");
const HttpClient = @import("../../HttpClient.zig");
const http = @import("../../../network/http.zig");
const Transfer = @import("../../../network/HttpClient.zig").Transfer;
const URL = @import("../../URL.zig");
const Mime = @import("../../Mime.zig");
@@ -47,7 +47,7 @@ _exec: *const Execution,
_proto: *XMLHttpRequestEventTarget,
_upload: ?*XMLHttpRequestUpload = null,
_arena: Allocator,
_http_response: ?HttpClient.Response = null,
_http_transfer: ?*Transfer = null,
// number of inflight requests, we can have multiple, e.g. xhr calling its own
// send from the onload callback
@@ -111,9 +111,9 @@ pub fn init(exec: *const Execution) !*XMLHttpRequest {
}
pub fn deinit(self: *XMLHttpRequest, page: *Page) void {
if (self._http_response) |resp| {
if (self._http_transfer) |resp| {
resp.abort(error.Abort);
self._http_response = null;
self._http_transfer = null;
}
if (self._on_ready_state_change) |func| {
@@ -182,9 +182,9 @@ pub fn setTimeout(self: *XMLHttpRequest, value: u32) void {
// TODO: url should be a union, as it can be multiple things
pub fn open(self: *XMLHttpRequest, method_: []const u8, url: [:0]const u8) !void {
// Abort any in-progress request
if (self._http_response) |transfer| {
if (self._http_transfer) |transfer| {
transfer.abort(error.Abort);
self._http_response = null;
self._http_transfer = null;
}
self._send_flag = false;
@@ -263,7 +263,7 @@ pub fn send(self: *XMLHttpRequest, body_: ?BodyInit, exec_: *const Execution) !v
self._active_requests += 1;
self._send_flag = true;
exec.makeRequest(.{
const transfer = exec.newRequest(.{
.ctx = self,
.url = self._url,
.method = self._method,
@@ -276,7 +276,6 @@ pub fn send(self: *XMLHttpRequest, body_: ?BodyInit, exec_: *const Execution) !v
.resource_type = .xhr,
.timeout_ms = self._timeout,
.notification = session.notification,
.start_callback = httpStartCallback,
.header_callback = httpHeaderDoneCallback,
.data_callback = httpDataCallback,
.done_callback = httpDoneCallback,
@@ -288,6 +287,23 @@ pub fn send(self: *XMLHttpRequest, body_: ?BodyInit, exec_: *const Execution) !v
self._send_flag = false;
return err;
};
// Held for abort() / open() / deinit; the error, shutdown and done
// callbacks clear it.
self._http_transfer = transfer;
if (comptime IS_DEBUG) {
log.debug(.http, "request start", .{ .method = self._method, .url = self._url, .source = "xhr" });
}
transfer.submit() catch |err| {
// If the failure aborted the transfer, httpErrorCallback already
// released the self-ref and cleared _http_transfer, making this
// release a guarded no-op.
self.releaseSelfRef();
self._send_flag = false;
return err;
};
}
// https://xhr.spec.whatwg.org/#the-upload-attribute
@@ -453,32 +469,24 @@ pub fn getResponseXML(self: *XMLHttpRequest, exec: *const Execution) !?*Node.Doc
}
}
fn httpStartCallback(response: HttpClient.Response) !void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(response.ctx));
if (comptime IS_DEBUG) {
log.debug(.http, "request start", .{ .method = self._method, .url = self._url, .source = "xhr" });
}
self._http_response = response;
}
fn httpHeaderCallback(response: HttpClient.Response, header: http.Header) !void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(response.ctx));
fn httpHeaderCallback(transfer: *Transfer, header: http.Header) !void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(transfer.req.ctx));
const joined = try std.fmt.allocPrint(self._arena, "{s}: {s}", .{ header.name, header.value });
try self._response_headers.append(self._arena, joined);
}
fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResult {
const self: *XMLHttpRequest = @ptrCast(@alignCast(response.ctx));
fn httpHeaderDoneCallback(transfer: *Transfer) !Transfer.HeaderResult {
const self: *XMLHttpRequest = @ptrCast(@alignCast(transfer.req.ctx));
if (comptime IS_DEBUG) {
log.debug(.http, "request header", .{
.source = "xhr",
.url = self._url,
.status = response.status(),
.status = transfer.responseStatus(),
});
}
if (response.contentType()) |ct| {
if (transfer.contentType()) |ct| {
self._response_mime = Mime.parse(ct) catch |e| {
log.info(.http, "invalid content type", .{
.content_Type = ct,
@@ -489,18 +497,18 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
};
}
var it = response.headerIterator();
var it = transfer.responseHeaderIterator();
while (it.next()) |hdr| {
const joined = try std.fmt.allocPrint(self._arena, "{s}: {s}", .{ hdr.name, hdr.value });
try self._response_headers.append(self._arena, joined);
}
self._response_status = response.status().?;
if (response.contentLength()) |cl| {
self._response_status = transfer.responseStatus().?;
if (transfer.getContentLength()) |cl| {
self._response_len = cl;
try self._response_data.ensureTotalCapacity(self._arena, cl);
}
self._response_url = try self._arena.dupeZ(u8, response.url());
self._response_url = try self._arena.dupeZ(u8, transfer.req.url);
const exec = self._exec;
@@ -515,8 +523,8 @@ fn httpHeaderDoneCallback(response: HttpClient.Response) !HttpClient.HeaderResul
return .proceed;
}
fn httpDataCallback(response: HttpClient.Response, data: []const u8) !void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(response.ctx));
fn httpDataCallback(transfer: *Transfer, data: []const u8) !void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(transfer.req.ctx));
try self._response_data.appendSlice(self._arena, data);
try self._proto.dispatch(.progress, .{
@@ -537,7 +545,7 @@ fn httpDoneCallback(ctx: *anyopaque) !void {
// Not that the request is done, the http/client will free the transfer
// object. It isn't safe to keep it around.
self._http_response = null;
self._http_transfer = null;
const exec = self._exec;
@@ -560,22 +568,22 @@ fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(ctx));
// http client will close it after an error, it isn't safe to keep around
self.handleError(err);
if (self._http_response != null) {
self._http_response = null;
if (self._http_transfer != null) {
self._http_transfer = null;
}
self.releaseSelfRef();
}
fn httpShutdownCallback(ctx: *anyopaque) void {
const self: *XMLHttpRequest = @ptrCast(@alignCast(ctx));
self._http_response = null;
self._http_transfer = null;
self.releaseSelfRef();
}
pub fn abort(self: *XMLHttpRequest) void {
self.handleError(error.Abort);
if (self._http_response) |resp| {
self._http_response = null;
if (self._http_transfer) |resp| {
self._http_transfer = null;
resp.abort(error.Abort);
}
self.releaseSelfRef();
+6 -4
View File
@@ -21,9 +21,12 @@ const lp = @import("lightpanda");
const App = @import("../App.zig");
const Inbox = @import("../Inbox.zig");
const Network = @import("../network/Network.zig");
const Notification = @import("../Notification.zig");
const WS = @import("../network/WS.zig");
const Network = @import("../network/Network.zig");
const Transfer = @import("../network/HttpClient.zig").Transfer;
const js = @import("../browser/js/js.zig");
const Browser = @import("../browser/Browser.zig");
const Session = @import("../browser/Session.zig");
@@ -32,7 +35,6 @@ const Page = @import("../browser/Page.zig");
const Mime = @import("../browser/Mime.zig");
const Element = @import("../browser/webapi/Element.zig");
const Label = @import("../browser/webapi/element/html/Label.zig");
const Transfer = @import("../browser/HttpClient.zig").Transfer;
const Connection = @import("Connection.zig");
const Incrementing = @import("id.zig").Incrementing;
@@ -981,8 +983,7 @@ pub const BrowserContext = struct {
// Encode the data in base64 by default, but don't encode
// for well known content-type.
.must_encode = blk: {
const response = msg.response;
if (response.contentType()) |ct| {
if (msg.transfer.contentType()) |ct| {
const mime = try Mime.parse(ct);
if (!mime.isText()) {
@@ -1458,5 +1459,6 @@ test "cdp: syncRequest short-circuits after disconnect" {
.cookie_origin = "",
.resource_type = .fetch,
.notification = undefined,
.shutdown_callback = @import("../network/HttpClient.zig").noopShutdown,
}));
}
+3 -3
View File
@@ -290,7 +290,7 @@ fn continueRequest(cmd: *CDP.Command) !void {
request.body = body;
}
try client.interception_layer.continueRequest(transfer);
try client.continueIntercepted(transfer);
return cmd.sendResult(null, .{});
}
@@ -396,7 +396,7 @@ fn fulfillRequest(cmd: *CDP.Command) !void {
body = buf;
}
try client.interception_layer.fulfillRequest(transfer, params.responseCode, params.responseHeaders orelse &.{}, body);
try client.fulfillIntercepted(transfer, params.responseCode, params.responseHeaders orelse &.{}, body);
return cmd.sendResult(null, .{});
}
@@ -417,7 +417,7 @@ fn failRequest(cmd: *CDP.Command) !void {
return cmd.sendResult(null, .{});
};
defer client.interception_layer.abortRequest(transfer);
defer transfer.abortParked(error.Abort);
log.info(.cdp, "request intercept", .{
.state = "fail",
+15 -41
View File
@@ -27,9 +27,9 @@ const URL = @import("../../browser/URL.zig");
const Mime = @import("../../browser/Mime.zig");
const Notification = @import("../../Notification.zig");
const timestamp = @import("../../datetime.zig").timestamp;
const Headers = @import("../../browser/HttpClient.zig").Headers;
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const Response = @import("../../browser/HttpClient.zig").Response;
const Headers = @import("../../network/HttpClient.zig").Headers;
const Transfer = @import("../../network/HttpClient.zig").Transfer;
const CdpStorage = @import("storage.zig");
@@ -91,7 +91,7 @@ fn setCacheDisabled(cmd: *CDP.Command) !void {
const bc = cmd.browser_context orelse return error.BrowserContextNotLoaded;
const client = &bc.cdp.browser.http_client;
client.cache_layer.disabled = params.cacheDisabled;
client.disableCache(params.cacheDisabled);
return cmd.sendResult(null, .{});
}
@@ -362,7 +362,7 @@ pub fn httpResponseHeaderDone(arena: Allocator, bc: *CDP.BrowserContext, msg: *c
.frameId = &id.toFrameId(req.frame_id),
.requestId = &id.toRequestId(transfer),
.loaderId = &id.toLoaderId(req.loader_id),
.response = ResponseWriter.init(arena, msg.response),
.response = ResponseWriter.init(arena, msg.transfer),
.hasExtraInfo = false, // TODO change after adding Network.responseReceivedExtraInfo
}, .{ .session_id = session_id });
}
@@ -447,12 +447,12 @@ pub const RequestWriter = struct {
const ResponseWriter = struct {
arena: Allocator,
response: *const Response,
transfer: *Transfer,
fn init(arena: Allocator, response: *const Response) ResponseWriter {
fn init(arena: Allocator, transfer: *Transfer) ResponseWriter {
return .{
.arena = arena,
.response = response,
.transfer = transfer,
};
}
@@ -461,15 +461,15 @@ const ResponseWriter = struct {
}
fn _jsonStringify(self: *const ResponseWriter, jws: anytype) !void {
const response = self.response;
const transfer = self.transfer;
try jws.beginObject();
{
try jws.objectField("url");
try jws.write(response.url());
try jws.write(transfer.req.url);
}
if (response.status()) |status| {
if (transfer.responseStatus()) |status| {
try jws.objectField("status");
try jws.write(status);
@@ -479,7 +479,7 @@ const ResponseWriter = struct {
{
const mime: Mime = blk: {
if (response.contentType()) |ct| {
if (transfer.contentType()) |ct| {
break :blk try Mime.parse(ct);
}
break :blk .unknown;
@@ -493,7 +493,7 @@ const ResponseWriter = struct {
{
try jws.objectField("fromDiskCache");
try jws.write(response.inner == .cached);
try jws.write(transfer._from_cache);
}
{
@@ -521,7 +521,7 @@ const ResponseWriter = struct {
// common to get these from a server (e.g. for Cache-Control), but
// Chrome joins these. So we have to too.
const arena = self.arena;
var it = response.headerIterator();
var it = transfer.responseHeaderIterator();
var map: std.StringArrayHashMapUnmanaged([]const u8) = .empty;
while (it.next()) |hdr| {
const gop = try map.getOrPut(arena, hdr.name);
@@ -866,7 +866,7 @@ test "cdp.Network: canClearBrowserCache" {
try ctx.expectSentResult(.{ .result = false }, .{ .id = 1 });
}
test "cdp.Network: setCacheDisabled disables cache" {
test "cdp.Network: setCacheDisabled" {
var ctx = try testing.context();
defer ctx.deinit();
_ = try ctx.loadBrowserContext(.{ .id = "BID-CD1" });
@@ -877,30 +877,4 @@ test "cdp.Network: setCacheDisabled disables cache" {
.params = .{ .cacheDisabled = true },
});
try ctx.expectSentResult(null, .{ .id = 1 });
const client = ctx.cdp().browser.http_client;
try testing.expectEqual(true, client.cache_layer.disabled);
}
test "cdp.Network: setCacheDisabled re-enables cache" {
var ctx = try testing.context();
defer ctx.deinit();
_ = try ctx.loadBrowserContext(.{ .id = "BID-CD2" });
try ctx.processMessage(.{
.id = 1,
.method = "Network.setCacheDisabled",
.params = .{ .cacheDisabled = true },
});
try ctx.expectSentResult(null, .{ .id = 1 });
try ctx.processMessage(.{
.id = 2,
.method = "Network.setCacheDisabled",
.params = .{ .cacheDisabled = false },
});
try ctx.expectSentResult(null, .{ .id = 2 });
const client = ctx.cdp().browser.http_client;
try testing.expectEqual(false, client.cache_layer.disabled);
}
+2 -2
View File
@@ -149,8 +149,8 @@ fn setLifecycleEventsEnabled(cmd: *CDP.Command) !void {
const http_client = frame._session.browser.http_client;
const http_active = http_client.http_active;
const http_next_tick = http_client.next_tick_count;
const total_network_activity = http_active + http_next_tick + http_client.interception_layer.intercepted;
const http_buffered = http_client.dispatch_count;
const total_network_activity = http_active + http_buffered + http_client.intercepted;
if (frame._notified_network_almost_idle.check(total_network_activity <= 2)) {
try sendPageLifecycle(bc, "networkAlmostIdle", now, frame_id, loader_id);
}
+1 -1
View File
@@ -40,7 +40,7 @@ pub fn toLoaderId(id: u32) [14]u8 {
// requestId has special requirements. If it's the main document navigation,
// then it should match the loader id.
const Transfer = @import("../browser/HttpClient.zig").Transfer;
const Transfer = @import("../network/HttpClient.zig").Transfer;
pub fn toRequestId(transfer: *const Transfer) [14]u8 {
if (transfer.req.resource_type == .document) {
return toLoaderId(transfer.req.loader_id);
+1 -1
View File
@@ -43,7 +43,7 @@ pub const forms = @import("browser/forms.zig");
pub const actions = @import("browser/actions.zig");
pub const structured_data = @import("browser/structured_data.zig");
pub const tools = @import("browser/tools.zig");
pub const HttpClient = @import("browser/HttpClient.zig");
pub const HttpClient = @import("network/HttpClient.zig");
pub const mcp = @import("mcp.zig");
pub const Agent = @import("agent/Agent.zig");
File diff suppressed because it is too large. Load diff
+232
View File
@@ -0,0 +1,232 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
// robots.txt gate for the HttpClient request pipeline. Answers allow/deny
// from the robot store; on a store miss it parks the transfer, coalesces
// concurrent requests for the same robots.txt behind a single internal
// fetch, and resumes (or fails) the parked transfers when it resolves.
const std = @import("std");
const lp = @import("lightpanda");
const URL = @import("../browser/URL.zig");
const Robots = @import("Robots.zig");
const Network = @import("Network.zig");
const Transfer = @import("HttpClient.zig").Transfer;
const log = lp.log;
const Allocator = std.mem.Allocator;
const RobotsGate = @This();
network: *Network,
allocator: Allocator,
pending: std.StringHashMapUnmanaged(std.ArrayList(*Transfer)) = .empty,
pub const Result = enum { allowed, blocked, pending };
pub fn deinit(self: *RobotsGate) void {
var it = self.pending.iterator();
while (it.next()) |entry| {
entry.value_ptr.deinit(self.allocator);
}
self.pending.deinit(self.allocator);
}
pub fn check(self: *RobotsGate, transfer: *Transfer) !Result {
const url = transfer.req.url;
const robots_url = try URL.getRobotsUrl(transfer.arena, url);
if (self.network.robot_store.get(robots_url)) |robot_entry| {
switch (robot_entry) {
.absent => return .allowed,
.present => |robots| {
if (robots.isAllowed(URL.getPathname(url))) {
return .allowed;
}
log.warn(.http, "blocked by robots", .{ .url = url });
return .blocked;
},
}
unreachable;
}
try self.fetchThenResume(robots_url, transfer);
return .pending;
}
fn fetchThenResume(self: *RobotsGate, robots_url: [:0]const u8, transfer: *Transfer) !void {
const entry = try self.pending.getOrPut(self.allocator, robots_url);
if (entry.found_existing == false) {
// A fetch for this robots.txt is already in flight, queue behind it.
try entry.value_ptr.append(self.allocator, transfer);
transfer.park(.robots);
return;
}
errdefer _ = self.pending.remove(robots_url);
entry.value_ptr.* = .empty;
try entry.value_ptr.append(self.allocator, transfer);
transfer.park(.robots);
errdefer {
entry.value_ptr.deinit(self.allocator);
transfer.unpark();
}
const robots_ctx = try transfer.arena.create(RobotsContext);
robots_ctx.* = .{
.gate = self,
.buffer = .empty,
.arena = transfer.arena,
.robots_url = robots_url,
};
log.debug(.browser, "fetching robots.txt", .{ .robots_url = robots_url });
// Only the parent's frame/loader ids (CDP correlation) and notification
// carry over — no cookies, credentials, headers, or timeout.
try transfer.client.request(.{
.url = robots_url,
.method = .GET,
.internal = true,
.resource_type = .fetch,
.frame_id = transfer.req.frame_id,
.loader_id = transfer.req.loader_id,
.notification = transfer.req.notification,
.cookie_jar = null,
.cookie_origin = robots_url,
.ctx = robots_ctx,
.header_callback = RobotsContext.headerCallback,
.data_callback = RobotsContext.dataCallback,
.done_callback = RobotsContext.doneCallback,
.error_callback = RobotsContext.errorCallback,
.shutdown_callback = RobotsContext.shutdownCallback,
}, transfer.owner);
}
fn flushPending(self: *RobotsGate, robots_url: [:0]const u8, allowed: bool) void {
var queued = self.pending.fetchRemove(robots_url) orelse return;
defer queued.value.deinit(self.allocator);
for (queued.value.items) |transfer| {
transfer.unpark();
if (!allowed) {
log.warn(.http, "blocked by robots", .{ .url = transfer.req.url });
transfer.failAsync(error.RobotsBlocked);
continue;
}
// Hand back to the pipeline; the robots gate is the last step
// before the network. If it fails before committing, clean up here.
transfer.client.resumeAfterRobots(transfer) catch |e| {
if (transfer.state == .created) {
transfer.abort(e);
}
};
}
}
// shutdown_callback is only called on owner shutdown. And if this robot's fetch
// is being shutdown, than any transfers waiting for it will be shutdown too.
fn flushPendingShutdown(self: *RobotsGate, robots_url: [:0]const u8) void {
var pending = self.pending.fetchRemove(robots_url) orelse return;
pending.value.deinit(self.allocator);
}
const RobotsContext = struct {
gate: *RobotsGate,
arena: Allocator,
robots_url: [:0]const u8,
buffer: std.ArrayList(u8),
status: u16 = 0,
fn headerCallback(transfer: *Transfer) anyerror!Transfer.HeaderResult {
const self: *RobotsContext = @ptrCast(@alignCast(transfer.req.ctx));
if (transfer.res.header) |hdr| {
log.debug(.browser, "robots status", .{ .status = hdr.status, .robots_url = self.robots_url });
self.status = hdr.status;
}
if (transfer.getContentLength()) |cl| {
try self.buffer.ensureTotalCapacity(self.arena, cl);
}
return .proceed;
}
fn dataCallback(transfer: *Transfer, data: []const u8) anyerror!void {
const self: *RobotsContext = @ptrCast(@alignCast(transfer.req.ctx));
try self.buffer.appendSlice(self.arena, data);
}
fn doneCallback(ctx_ptr: *anyopaque) anyerror!void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
const gate = self.gate;
const robots_url = self.robots_url;
var allowed = true;
const network = gate.network;
switch (self.status) {
200 => {
if (self.buffer.items.len > 0) {
const robots: ?Robots = network.robot_store.robotsFromBytes(
network.config.http_headers.user_agent,
self.buffer.items,
) catch blk: {
log.warn(.browser, "failed to parse robots", .{ .robots_url = robots_url });
try network.robot_store.putAbsent(robots_url);
break :blk null;
};
if (robots) |r| {
try network.robot_store.put(robots_url, r);
const path = URL.getPathname(gate.pending.get(robots_url).?.items[0].req.url);
allowed = r.isAllowed(path);
}
}
},
404 => {
log.debug(.http, "robots not found", .{ .url = robots_url });
try network.robot_store.putAbsent(robots_url);
},
else => {
log.debug(.http, "unexpected status on robots", .{
.url = robots_url,
.status = self.status,
});
try network.robot_store.putAbsent(robots_url);
},
}
gate.flushPending(robots_url, allowed);
}
fn errorCallback(ctx_ptr: *anyopaque, err: anyerror) void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
log.warn(.http, "robots fetch failed", .{ .err = err });
self.gate.flushPending(self.robots_url, true);
}
fn shutdownCallback(ctx_ptr: *anyopaque) void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
log.debug(.http, "robots fetch shutdown", .{});
self.gate.flushPendingShutdown(self.robots_url);
}
};
+1 -6
View File
@@ -308,11 +308,6 @@ pub const ResponseHead = struct {
redirect_count: u32,
_content_type_len: usize = 0,
_content_type: [MAX_CONTENT_TYPE_LEN]u8 = undefined,
// this is normally an empty list, but if the response is being injected
// than it'll be populated. It isn't meant to be used directly, but should
// be used through the transfer.responseHeaderIterator() which abstracts
// whether the headers are from a live curl easy handle, or injected.
_injected_headers: []const Header = &.{},
pub fn contentType(self: *ResponseHead) ?[]u8 {
if (self._content_type_len == 0) {
@@ -364,7 +359,7 @@ pub const Connection = struct {
pub const Transport = union(enum) {
none, // used for cases that manage their own connection, e.g. telemetry
http: *@import("../browser/HttpClient.zig").Transfer,
http: *@import("HttpClient.zig").Transfer,
websocket: *@import("../browser/webapi/net/WebSocket.zig"),
};
-386
View File
@@ -1,386 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const lp = @import("lightpanda");
const Layer = @import("../../browser/HttpClient.zig").Layer;
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const Response = @import("../../browser/HttpClient.zig").Response;
const Cache = @import("../cache/Cache.zig");
const CachedMetadata = @import("../cache/Cache.zig").CachedMetadata;
const CachedResponse = @import("../cache/Cache.zig").CachedResponse;
const HeaderResult = @import("../../browser/HttpClient.zig").HeaderResult;
const Forward = @import("Forward.zig");
const log = lp.log;
const IS_DEBUG = @import("builtin").mode == .Debug;
const CacheLayer = @This();
next: Layer = undefined,
disabled: bool = false,
pub fn layer(self: *CacheLayer) Layer {
return .{
.ptr = self,
.vtable = &.{
.request = request,
},
};
}
fn request(ptr: *anyopaque, transfer: *Transfer) anyerror!void {
const self: *CacheLayer = @ptrCast(@alignCast(ptr));
const req = &transfer.req;
if (self.disabled or req.method != .GET) {
return self.next.request(transfer);
}
const arena = transfer.arena;
var iter = req.headers.iterator();
const req_header_list = try iter.collect(arena);
const cached = transfer.client.network.cache.?.get(arena, .{
.url = req.url,
.timestamp = std.time.timestamp(),
.request_headers = req_header_list.items,
}) orelse {
// Cache miss: install wrappers so we can inspect the response and decide
// whether to write the body into the cache when it's done.
try installCacheContext(arena, transfer, null);
return self.next.request(transfer);
};
if (!cached.expired) {
const ctx = try arena.create(CachedResponse);
ctx.* = cached;
try transfer.client.runNextTick(transfer, ctx, .{
.run = struct {
fn run(t: *Transfer, ctx_ptr: ?*anyopaque) void {
defer t.deinit();
const c: *CachedResponse = @ptrCast(@alignCast(ctx_ptr.?));
serveFromCache(t, c) catch |err| {
t.req.error_callback(t.req.ctx, err);
};
}
}.run,
.abort = struct {
fn abort(ctx_ptr: ?*anyopaque) void {
const c: *CachedResponse = @ptrCast(@alignCast(ctx_ptr.?));
switch (c.data) {
.buffer => |_| {},
.file => |f| f.file.close(),
}
}
}.abort,
});
return;
}
if (cached.metadata.hasValidators()) {
if (cached.metadata.etag) |etag| {
log.debug(.cache, "revalidate with etag", .{ .url = req.url, .etag = etag });
const header_value = try std.fmt.allocPrintSentinel(arena, "If-None-Match: {s}", .{etag}, 0);
try req.headers.add(header_value);
}
if (cached.metadata.last_modified) |lm| {
log.debug(.cache, "revalidate with last-modified", .{ .url = req.url, .last_modified = lm });
const header_value = try std.fmt.allocPrintSentinel(arena, "If-Modified-Since: {s}", .{lm}, 0);
try req.headers.add(header_value);
}
try installCacheContext(arena, transfer, cached);
} else {
defer cached.data.deinit();
// If it is expired w/o validators, evict from Cache.
transfer.client.network.cache.?.evict(req.url);
try installCacheContext(arena, transfer, null);
}
return self.next.request(transfer);
}
fn installCacheContext(
arena: std.mem.Allocator,
transfer: *Transfer,
stale_entry: ?CachedResponse,
) !void {
const req = &transfer.req;
const ctx = try arena.create(CacheContext);
ctx.* = .{
.arena = arena,
.transfer = transfer,
.forward = Forward.capture(req),
.req_url = req.url,
.req_headers = req.headers,
.stale_entry = stale_entry,
};
req.ctx = ctx;
req.header_callback = CacheContext.headerCallback;
req.data_callback = CacheContext.dataCallback;
req.done_callback = CacheContext.doneCallback;
req.error_callback = CacheContext.errorCallback;
if (ctx.forward.start != null) req.start_callback = CacheContext.startCallback;
if (ctx.forward.shutdown != null) req.shutdown_callback = CacheContext.shutdownCallback;
}
fn forwardFromCache(
transfer: *Transfer,
forward: *Forward,
cached: *const CachedResponse,
) !void {
transfer.req.notification.dispatch(
.http_request_served_from_cache,
&.{ .transfer = transfer },
);
const req = &transfer.req;
const response = Response.fromCached(req.ctx, cached);
defer cached.data.deinit();
try forward.forwardStart(response);
const result = try forward.forwardHeader(response);
if (result == .abort) return error.Abort;
switch (cached.data) {
.buffer => |data| {
if (data.len > 0) try forward.forwardData(response, data);
},
.file => |f| {
const file = f.file;
var buf: [1024]u8 = undefined;
var file_reader = file.reader(&buf);
try file_reader.seekTo(f.offset);
const reader = &file_reader.interface;
var read_buf: [1024]u8 = undefined;
var remaining = f.len;
while (remaining > 0) {
const read_len = @min(read_buf.len, remaining);
const n = try reader.readSliceShort(read_buf[0..read_len]);
if (n == 0) break;
remaining -= n;
try forward.forwardData(response, read_buf[0..n]);
}
},
}
try forward.forwardDone();
}
fn serveFromCache(transfer: *Transfer, cached: *const CachedResponse) !void {
transfer.req.notification.dispatch(
.http_request_served_from_cache,
&.{ .transfer = transfer },
);
const req = &transfer.req;
const response = Response.fromCached(req.ctx, cached);
defer cached.data.deinit();
if (req.start_callback) |cb| {
try cb(response);
}
const result = try req.header_callback(response);
if (result == .abort) {
return error.Abort;
}
switch (cached.data) {
.buffer => |data| {
if (data.len > 0) {
try req.data_callback(response, data);
}
},
.file => |f| {
const file = f.file;
var buf: [1024]u8 = undefined;
var file_reader = file.reader(&buf);
try file_reader.seekTo(f.offset);
const reader = &file_reader.interface;
var read_buf: [1024]u8 = undefined;
var remaining = f.len;
while (remaining > 0) {
const read_len = @min(read_buf.len, remaining);
const n = try reader.readSliceShort(read_buf[0..read_len]);
if (n == 0) break;
remaining -= n;
try req.data_callback(response, read_buf[0..n]);
}
},
}
try req.done_callback(req.ctx);
}
const CacheContext = struct {
arena: std.mem.Allocator,
transfer: *Transfer,
forward: Forward,
req_url: [:0]const u8,
req_headers: @import("../http.zig").Headers,
pending_metadata: ?*CachedMetadata = null,
stale_entry: ?CachedResponse = null,
fn startCallback(response: Response) anyerror!void {
const self: *CacheContext = @ptrCast(@alignCast(response.ctx));
return self.forward.forwardStart(response);
}
fn dataCallback(response: Response, chunk: []const u8) anyerror!void {
const self: *CacheContext = @ptrCast(@alignCast(response.ctx));
return self.forward.forwardData(response, chunk);
}
fn headerCallback(response: Response) anyerror!HeaderResult {
const self: *CacheContext = @ptrCast(@alignCast(response.ctx));
// For non-transfer responses (fulfilled by interception, or future
// cached-while-cached cases), there's nothing to inspect for caching
// decisions — just forward.
const transfer = switch (response.inner) {
.transfer => |t| t,
else => return self.forward.forwardHeader(response),
};
const arena = self.arena;
const conn = transfer._conn.?;
var rh = &transfer.res.header.?;
if (self.stale_entry != null and rh.status == 304) {
const stale = self.stale_entry.?;
self.stale_entry = null;
var iter = response.headerIterator();
const headers = try iter.collect(arena);
transfer.client.network.cache.?.renew(
arena,
.{
.url = self.req_url,
.timestamp = std.time.timestamp(),
.headers = headers.items,
},
) catch |err| {
log.warn(.cache, "renew failed", .{ .err = err });
};
try forwardFromCache(transfer, &self.forward, &stale);
return .handled;
}
if (self.stale_entry) |stale| {
stale.data.deinit();
self.stale_entry = null;
}
const vary = if (conn.getResponseHeader("vary", 0)) |h| h.value else null;
const maybe_cm = try Cache.tryCache(
arena,
std.time.timestamp(),
self.req_url,
rh.status,
rh.contentType(),
if (conn.getResponseHeader("cache-control", 0)) |h| h.value else null,
vary,
if (conn.getResponseHeader("age", 0)) |h| h.value else null,
if (conn.getResponseHeader("etag", 0)) |h| h.value else null,
if (conn.getResponseHeader("last-modified", 0)) |h| h.value else null,
conn.getResponseHeader("set-cookie", 0) != null,
conn.getResponseHeader("authorization", 0) != null,
);
if (maybe_cm) |cm| {
var iter = transfer.responseHeaderIterator();
var header_list = try iter.collect(arena);
const end_of_response = header_list.items.len;
if (vary) |vary_str| {
var req_it = self.req_headers.iterator();
while (req_it.next()) |hdr| {
var vary_iter = std.mem.splitScalar(u8, vary_str, ',');
while (vary_iter.next()) |part| {
const name = std.mem.trim(u8, part, &std.ascii.whitespace);
if (std.ascii.eqlIgnoreCase(hdr.name, name)) {
try header_list.append(arena, .{
.name = try arena.dupe(u8, hdr.name),
.value = try arena.dupe(u8, hdr.value),
});
}
}
}
}
const metadata = try arena.create(CachedMetadata);
metadata.* = cm;
metadata.headers = header_list.items[0..end_of_response];
metadata.vary_headers = header_list.items[end_of_response..];
self.pending_metadata = metadata;
}
return self.forward.forwardHeader(response);
}
fn doneCallback(ctx: *anyopaque) anyerror!void {
const self: *CacheContext = @ptrCast(@alignCast(ctx));
const transfer = self.transfer;
if (self.pending_metadata) |metadata| {
const cache = &transfer.client.network.cache.?;
if (comptime IS_DEBUG) {
log.debug(.browser, "http cache", .{ .key = self.req_url, .metadata = metadata });
}
cache.put(metadata.*, transfer.res.stream_buffer.items) catch |err| {
log.warn(.http, "cache put failed", .{ .err = err });
};
}
return self.forward.forwardDone();
}
fn shutdownCallback(ctx: *anyopaque) void {
const self: *CacheContext = @ptrCast(@alignCast(ctx));
if (self.stale_entry) |entry| {
self.stale_entry = null;
entry.data.deinit();
}
self.forward.forwardShutdown();
}
fn errorCallback(ctx: *anyopaque, e: anyerror) void {
const self: *CacheContext = @ptrCast(@alignCast(ctx));
if (self.stale_entry) |entry| {
self.stale_entry = null;
entry.data.deinit();
}
self.forward.forwardErr(e);
}
};
-361
View File
@@ -1,361 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const lp = @import("lightpanda");
const log = lp.log;
const Network = @import("../Network.zig");
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const Response = @import("../../browser/HttpClient.zig").Response;
const Layer = @import("../../browser/HttpClient.zig").Layer;
const StableResponse = @import("../../browser/HttpClient.zig").StableResponse;
const Forward = @import("Forward.zig");
const HeaderResult = @import("../../browser/HttpClient.zig").HeaderResult;
const DeferringLayer = @This();
allocator: std.mem.Allocator,
network: *Network,
next: Layer = undefined,
active: std.DoublyLinkedList = .{},
pub fn layer(self: *DeferringLayer) Layer {
return .{
.ptr = self,
.vtable = &.{ .request = request },
};
}
pub fn deinit(self: *DeferringLayer) void {
self.drainAll();
}
fn request(ptr: *anyopaque, transfer: *Transfer) anyerror!void {
const self: *DeferringLayer = @ptrCast(@alignCast(ptr));
if (transfer.req.internal) {
return self.next.request(transfer);
}
const arena = try self.network.app.arena_pool.acquire(.small, "DeferringContext");
errdefer self.network.app.arena_pool.release(arena);
// this might outlive the transfer, we need to dupe eveyrthing we'll need to use
const ctx = try arena.create(DeferredContext);
ctx.* = .{
.arena = arena,
.layer = self,
.transfer = transfer,
.frame_id = transfer.req.frame_id,
.url = try arena.dupeZ(u8, transfer.req.url),
.forward = Forward.capture(&transfer.req),
};
self.active.append(&ctx.node);
errdefer self.active.remove(&ctx.node);
transfer.req.ctx = ctx;
transfer.req.start_callback = if (ctx.forward.start != null) DeferredContext.startCallback else null;
transfer.req.header_callback = DeferredContext.headerCallback;
transfer.req.data_callback = DeferredContext.dataCallback;
transfer.req.done_callback = DeferredContext.doneCallback;
transfer.req.error_callback = DeferredContext.errorCallback;
transfer.req.shutdown_callback = if (ctx.forward.shutdown != null) DeferredContext.shutdownCallback else null;
return self.next.request(transfer);
}
pub fn flushFrame(self: *DeferringLayer, frame_id: u32) void {
// DeferredContext.fire() can re-enter flushFrame, so we'll capture
// ready items in this list, so that a reentrant flushFrame doesn't mutate
// self.active while we're iterating.
var ready: std.DoublyLinkedList = .{};
var node = self.active.first;
while (node) |n| {
node = n.next;
const ctx: *DeferredContext = @fieldParentPtr("node", n);
if (!ctx.deferring) {
continue;
}
// captured frame_id, not ctx.transfer: the transfer may be freed.
if (ctx.frame_id != frame_id) {
continue;
}
if (ctx.terminal) {
self.active.remove(n);
ready.append(n);
} else {
ctx.firePartial();
ctx.deferring = false;
}
}
// ready is local, ctx.fire() re-entering flushFrame can't invalidate it.
while (ready.popFirst()) |n| {
const ctx: *DeferredContext = @fieldParentPtr("node", n);
ctx.fire();
}
}
/// Drop orphaned deferred contexts for a frame that's going away. A `terminal`
/// context's transfer already completed while deferred, so it's been deinited
/// and unlinked from the owner — abortOwner can't reach it, yet it lingers in
/// `active` pointing at a forward target (the Fetch) whose arena page teardown
/// is about to free, and a later flushFrame would fire into it. Non-terminal
/// contexts still have a live transfer that cleans them up itself.
pub fn cancelFrame(self: *DeferringLayer, frame_id: u32) void {
var node = self.active.first;
while (node) |n| {
node = n.next;
const ctx: *DeferredContext = @fieldParentPtr("node", n);
if (ctx.frame_id != frame_id or !ctx.terminal) {
continue;
}
self.active.remove(n);
ctx.deinit();
}
}
pub fn drainAll(self: *DeferringLayer) void {
while (self.active.popFirst()) |node| {
const ctx: *DeferredContext = @fieldParentPtr("node", node);
ctx.deinit();
}
}
const DeferredContext = struct {
arena: std.mem.Allocator,
layer: *DeferringLayer,
transfer: *Transfer,
frame_id: u32,
url: [:0]const u8,
forward: Forward,
node: std.DoublyLinkedList.Node = .{},
buffered: std.ArrayList(BufferedEvent) = .{},
done: bool = false,
deferring: bool = false,
terminal: bool = false,
stable_resp: ?StableResponse = null,
const BufferedEvent = union(enum) {
start,
header,
data: []const u8,
done,
err: anyerror,
};
fn deinit(self: *DeferredContext) void {
self.layer.network.app.arena_pool.release(self.arena);
}
fn setStableResponse(self: *DeferredContext, response: Response) !void {
if (self.stable_resp == null) {
self.stable_resp = try Response.toStable(response, self.arena);
}
}
fn shouldDefer(self: *DeferredContext) bool {
const req = self.transfer.req;
const blocking_id = self.transfer.client.blocking_requests.get(req.frame_id) orelse return false;
return self.transfer.id != blocking_id;
}
fn startCallback(response: Response) anyerror!void {
const self: *DeferredContext = @ptrCast(@alignCast(response.ctx));
if (!self.deferring and !self.shouldDefer()) {
return self.forward.forwardStart(response);
}
log.debug(.http, "deferring start callback", .{ .url = self.url });
try self.setStableResponse(response);
self.deferring = true;
try self.buffered.append(self.arena, .start);
}
fn headerCallback(response: Response) anyerror!HeaderResult {
const self: *DeferredContext = @ptrCast(@alignCast(response.ctx));
if (!self.deferring and !self.shouldDefer()) {
return self.forward.forwardHeader(response);
}
log.debug(.http, "deferring header callback", .{ .url = self.url });
try self.setStableResponse(response);
self.deferring = true;
try self.buffered.append(self.arena, .header);
return .proceed;
}
fn dataCallback(response: Response, chunk: []const u8) anyerror!void {
const self: *DeferredContext = @ptrCast(@alignCast(response.ctx));
if (!self.deferring and !self.shouldDefer()) {
return self.forward.forwardData(response, chunk);
}
log.debug(.http, "deferring data callback", .{ .url = self.url });
try self.setStableResponse(response);
self.deferring = true;
try self.buffered.append(self.arena, .{ .data = try self.arena.dupe(u8, chunk) });
}
fn doneCallback(ctx: *anyopaque) anyerror!void {
const self: *DeferredContext = @ptrCast(@alignCast(ctx));
if (!self.deferring and !self.shouldDefer()) {
defer self.deinit();
self.done = true;
self.layer.active.remove(&self.node);
return self.forward.forwardDone();
}
log.debug(.http, "deferring done callback", .{ .url = self.url });
self.deferring = true;
self.terminal = true;
try self.buffered.append(self.arena, .done);
}
fn errorCallback(ctx: *anyopaque, err: anyerror) void {
const self: *DeferredContext = @ptrCast(@alignCast(ctx));
if (!self.deferring and !self.shouldDefer()) {
defer self.deinit();
self.done = true;
self.layer.active.remove(&self.node);
self.forward.forwardErr(err);
return;
}
log.debug(.http, "deferring error callback", .{ .url = self.url, .err = err });
self.deferring = true;
self.terminal = true;
self.buffered.append(self.arena, .{ .err = err }) catch {};
}
fn shutdownCallback(ctx: *anyopaque) void {
const self: *DeferredContext = @ptrCast(@alignCast(ctx));
if (self.done) return;
defer self.deinit();
self.done = true;
self.layer.active.remove(&self.node);
log.debug(.http, "deferring shutdown callback", .{});
self.forward.forwardShutdown();
}
fn fire(self: *DeferredContext) void {
defer self.deinit();
for (self.buffered.items) |event| {
switch (event) {
.start => {
const stable_response = self.stable_resp orelse @panic("stable_resp must be set for start events");
const response = Response.fromStable(&stable_response);
self.forward.forwardStart(response) catch |err| {
log.err(.http, "deferred start callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
},
.header => {
const stable_response = self.stable_resp orelse @panic("stable_resp must be set for header events");
const response = Response.fromStable(&stable_response);
const result = self.forward.forwardHeader(response) catch |err| {
log.err(.http, "deferred header callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
if (result == .abort) {
self.forward.forwardErr(error.Abort);
return;
}
},
.data => |chunk| {
const stable_response = self.stable_resp orelse @panic("stable_resp must be set for data events");
const response = Response.fromStable(&stable_response);
self.forward.forwardData(response, chunk) catch |err| {
log.err(.http, "deferred data callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
},
.done => {
self.forward.forwardDone() catch |err| {
log.err(.http, "deferred done callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
};
return;
},
.err => |err| {
self.forward.forwardErr(err);
return;
},
}
}
}
fn firePartial(self: *DeferredContext) void {
const stable_response = self.stable_resp orelse @panic("stable_resp must be set for any of the partial fire events");
const response = Response.fromStable(&stable_response);
for (self.buffered.items) |event| {
switch (event) {
.start => {
self.forward.forwardStart(response) catch |err| {
log.err(.http, "defer part start callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
},
.header => {
const result = self.forward.forwardHeader(response) catch |err| {
log.err(.http, "defer part header callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
if (result == .abort) {
self.forward.forwardErr(error.Abort);
return;
}
},
.data => |chunk| {
self.forward.forwardData(response, chunk) catch |err| {
log.err(.http, "defer part data callback", .{ .err = err, .url = self.url });
self.forward.forwardErr(err);
return;
};
},
.done, .err => @panic("firePartial cant fire terminal events"),
}
}
self.buffered.clearRetainingCapacity();
}
};
-76
View File
@@ -1,76 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
// A snapshot of the original ctx + callbacks from a Request, taken before a
// layer overwrites them with its own wrappers. The layer's wrapper callbacks
// call forwardX(...) to invoke the captured originals with the original ctx.
const Request = @import("../../browser/HttpClient.zig").Request;
const Response = @import("../../browser/HttpClient.zig").Response;
const HeaderResult = @import("../../browser/HttpClient.zig").HeaderResult;
const Forward = @This();
ctx: *anyopaque,
start: ?Request.StartCallback,
header: Request.HeaderCallback,
data: Request.DataCallback,
done: Request.DoneCallback,
err: Request.ErrorCallback,
shutdown: ?Request.ShutdownCallback,
pub fn capture(req: *const Request) Forward {
return .{
.ctx = req.ctx,
.start = req.start_callback,
.header = req.header_callback,
.data = req.data_callback,
.done = req.done_callback,
.err = req.error_callback,
.shutdown = req.shutdown_callback,
};
}
pub fn forwardStart(self: Forward, response: Response) anyerror!void {
var fwd = response;
fwd.ctx = self.ctx;
if (self.start) |cb| try cb(fwd);
}
pub fn forwardHeader(self: Forward, response: Response) anyerror!HeaderResult {
var fwd = response;
fwd.ctx = self.ctx;
return self.header(fwd);
}
pub fn forwardData(self: Forward, response: Response, chunk: []const u8) anyerror!void {
var fwd = response;
fwd.ctx = self.ctx;
return self.data(fwd, chunk);
}
pub fn forwardDone(self: Forward) anyerror!void {
return self.done(self.ctx);
}
pub fn forwardErr(self: Forward, e: anyerror) void {
self.err(self.ctx, e);
}
pub fn forwardShutdown(self: Forward) void {
if (self.shutdown) |cb| cb(self.ctx);
}
-327
View File
@@ -1,327 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const builtin = @import("builtin");
const lp = @import("lightpanda");
const log = lp.log;
const IS_DEBUG = builtin.mode == .Debug;
const http = @import("../http.zig");
const HttpClient = @import("../../browser/HttpClient.zig");
const Request = @import("../../browser/HttpClient.zig").Request;
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const Response = @import("../../browser/HttpClient.zig").Response;
const FulfilledResponse = @import("../../browser/HttpClient.zig").FulfilledResponse;
const Layer = @import("../../browser/HttpClient.zig").Layer;
const Forward = @import("Forward.zig");
const HeaderResult = @import("../../browser/HttpClient.zig").HeaderResult;
const InterceptionLayer = @This();
// Count of intercepted requests. The client doesn't track intercepted transfers
// on its own active counters: once intercepted, a transfer leaves the layer
// chain and waits for the interceptor (CDP) to call continue/abort/fulfill.
// We track them here so the network-idle / network-almost-idle CDP lifecycle
// events don't fire prematurely.
intercepted: usize = 0,
next: Layer = undefined,
pub fn layer(self: *InterceptionLayer) Layer {
return .{
.ptr = self,
.vtable = &.{ .request = request },
};
}
fn request(ptr: *anyopaque, transfer: *Transfer) anyerror!void {
const self: *InterceptionLayer = @ptrCast(@alignCast(ptr));
const req = &transfer.req;
const ctx = try transfer.arena.create(InterceptContext);
ctx.* = .{
.layer = self,
.transfer = transfer,
.forward = Forward.capture(req),
};
// Install our wrappers on the transfer's request. The interceptor wants to
// observe every callback (start/header/data/done/err/shutdown) so it can
// mirror the Network.* CDP events.
req.ctx = ctx;
if (ctx.forward.start != null) req.start_callback = InterceptContext.startCallback;
req.header_callback = InterceptContext.headerCallback;
req.data_callback = InterceptContext.dataCallback;
req.done_callback = InterceptContext.doneCallback;
req.error_callback = InterceptContext.errorCallback;
if (ctx.forward.shutdown != null) req.shutdown_callback = InterceptContext.shutdownCallback;
req.notification.dispatch(.http_request_start, &.{ .transfer = transfer });
var wait_for_interception = false;
req.notification.dispatch(.http_request_intercept, &.{
.transfer = transfer,
.wait_for_interception = &wait_for_interception,
});
log.debug(.http, "interception check", .{
.wait_for_interception = wait_for_interception,
.intercepted = self.intercepted,
.url = req.url,
});
if (!wait_for_interception) {
return self.next.request(transfer);
}
// Paused: the CDP listener stashed `transfer` and will eventually call
// continueRequest / abortRequest / fulfillRequest. Until then, CDP owns
// the transfer's lifecycle. Park keeps the outer Client.request errdefer
// from tearing it down.
self.intercepted += 1;
transfer.park(.intercept_request);
if (comptime IS_DEBUG) {
log.debug(.http, "wait for interception", .{ .intercepted = self.intercepted });
}
}
pub const InterceptContext = struct {
layer: *InterceptionLayer,
transfer: *Transfer,
forward: Forward,
content_length: usize = 0,
fn startCallback(response: Response) anyerror!void {
const self: *InterceptContext = @ptrCast(@alignCast(response.ctx));
log.debug(.http, "intercept start", .{ .url = self.transfer.req.url });
return self.forward.forwardStart(response);
}
fn headerCallback(response: Response) anyerror!HeaderResult {
const self: *InterceptContext = @ptrCast(@alignCast(response.ctx));
log.debug(.http, "intercept header", .{
.url = self.transfer.req.url,
.status = response.status(),
.content_length = response.contentLength(),
});
self.content_length = response.contentLength() orelse 0;
self.transfer.req.notification.dispatch(.http_response_header_done, &.{
.transfer = self.transfer,
.response = &response,
});
return self.forward.forwardHeader(response);
}
fn dataCallback(response: Response, chunk: []const u8) anyerror!void {
const self: *InterceptContext = @ptrCast(@alignCast(response.ctx));
log.debug(.http, "intercept data", .{
.url = self.transfer.req.url,
.len = chunk.len,
});
self.transfer.req.notification.dispatch(.http_response_data, &.{
.data = chunk,
.transfer = self.transfer,
});
return self.forward.forwardData(response, chunk);
}
fn doneCallback(ctx: *anyopaque) anyerror!void {
const self: *InterceptContext = @ptrCast(@alignCast(ctx));
log.debug(.http, "intercept done", .{
.url = self.transfer.req.url,
.content_length = self.content_length,
});
self.transfer.req.notification.dispatch(.http_request_done, &.{
.transfer = self.transfer,
.content_length = self.content_length,
});
return self.forward.forwardDone();
}
fn errorCallback(ctx: *anyopaque, err: anyerror) void {
const self: *InterceptContext = @ptrCast(@alignCast(ctx));
log.debug(.http, "intercept error", .{
.url = self.transfer.req.url,
.err = err,
});
self.transfer.req.notification.dispatch(.http_request_fail, &.{
.transfer = self.transfer,
.err = err,
});
self.forward.forwardErr(err);
}
fn shutdownCallback(ctx: *anyopaque) void {
const self: *InterceptContext = @ptrCast(@alignCast(ctx));
log.debug(.http, "intercept shutdown", .{ .url = self.transfer.req.url });
self.transfer.req.notification.dispatch(.http_request_fail, &.{
.transfer = self.transfer,
.err = error.Shutdown,
});
self.forward.forwardShutdown();
}
};
// CDP-driven resolution entry points. The transfer was paused inside `request`
// (state = .parked = .intercept_request). One of these three is called by CDP
// to resume / drop the transfer.
pub fn continueRequest(self: *InterceptionLayer, transfer: *Transfer) anyerror!void {
if (comptime IS_DEBUG) {
lp.assert(self.intercepted > 0, "InterceptionLayer.continueRequest", .{ .value = self.intercepted });
log.debug(.http, "continue transfer", .{ .intercepted = self.intercepted });
}
// Resume the layer chain. Ownership is re-handed to whichever subsequent
// layer commits the transfer (queue, multi, or another park). If the
// chain fails before any commit, we clean up here — mirror the errdefer
// pattern in Client.request.
transfer.unpark();
self.next.request(transfer) catch |err| {
if (transfer.state == .created) {
transfer.abort(err);
}
return err;
};
}
pub fn abortRequest(self: *InterceptionLayer, transfer: *Transfer) void {
if (comptime IS_DEBUG) {
lp.assert(self.intercepted > 0, "InterceptionLayer.abortRequest", .{ .value = self.intercepted });
log.debug(.http, "abort transfer", .{ .intercepted = self.intercepted });
}
transfer.abortParked(error.Abort);
}
pub fn fulfillRequest(
self: *InterceptionLayer,
transfer: *Transfer,
status: u16,
headers: []const http.Header,
body: ?[]const u8,
) !void {
if (comptime IS_DEBUG) {
lp.assert(self.intercepted > 0, "InterceptionLayer.fulfillRequest", .{ .value = self.intercepted });
log.debug(.http, "fulfill transfer", .{ .intercepted = self.intercepted });
}
// Leave the parked state (accounting `intercepted` exactly once) and move to
// .completing BEFORE running the user callbacks.
transfer.unpark();
if (HttpClient.isRedirectStatus(status)) {
if (findLocation(headers)) |location| {
fulfilledRedirect(transfer, status, headers, location) catch |err| {
if (transfer.state == .created) {
transfer.abort(err);
}
return err;
};
self.next.request(transfer) catch |err| {
if (transfer.state == .created) {
transfer.abort(err);
}
return err;
};
return;
}
}
// Not a redirect: move to .completing BEFORE running the user callbacks.
transfer.state = .completing;
defer transfer.deinit();
// `done` flips true once we've called the user's done_callback. If
// done_callback itself throws, the user already saw their end-of-flow
// notification; suppress error_callback to avoid double-notify.
var done: bool = false;
fulfillInner(&transfer.req, status, headers, body, &done) catch |err| {
if (!done) {
// safe here despite the defer transfer.deinit() above since the
// state == .completing
transfer.abort(err);
}
return err;
};
}
fn fulfillInner(
req: *Request,
status: u16,
headers: []const http.Header,
body: ?[]const u8,
done: *bool,
) !void {
const fulfilled = FulfilledResponse{
.status = status,
.url = req.url,
.headers = headers,
.body = body,
};
const response = Response.fromFulfilled(req.ctx, &fulfilled);
if (req.start_callback) |cb| {
try cb(response);
}
const result = try req.header_callback(response);
if (result == .abort) {
return error.Abort;
}
if (body) |b| {
try req.data_callback(response, b);
}
done.* = true;
try req.done_callback(req.ctx);
}
fn fulfilledRedirect(transfer: *Transfer, status: u16, headers: []const http.Header, location: []const u8) !void {
// retrieve cookies from the fulfilled response's headers.
if (transfer.req.cookie_jar) |jar| {
for (headers) |hdr| {
if (std.ascii.eqlIgnoreCase(hdr.name, "set-cookie")) {
try jar.populateFromResponse(transfer.req.url, hdr.value);
}
}
}
try transfer.applyRedirectTarget(transfer.req.url, location, status);
}
fn findLocation(headers: []const http.Header) ?[]const u8 {
for (headers) |hdr| {
if (std.ascii.eqlIgnoreCase(hdr.name, "location")) {
return hdr.value;
}
}
return null;
}
-291
View File
@@ -1,291 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const lp = @import("lightpanda");
const URL = @import("../../browser/URL.zig");
const Layer = @import("../../browser/HttpClient.zig").Layer;
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const Response = @import("../../browser/HttpClient.zig").Response;
const HeaderResult = @import("../../browser/HttpClient.zig").HeaderResult;
const Robots = @import("../Robots.zig");
const Network = @import("../Network.zig");
const log = lp.log;
const Allocator = std.mem.Allocator;
const RobotsLayer = @This();
next: Layer = undefined,
network: *Network,
allocator: Allocator,
pending: std.StringHashMapUnmanaged(std.ArrayList(*Transfer)) = .empty,
pub fn layer(self: *RobotsLayer) Layer {
return .{
.ptr = self,
.vtable = &.{
.request = request,
},
};
}
pub fn deinit(self: *RobotsLayer, allocator: Allocator) void {
var it = self.pending.iterator();
while (it.next()) |entry| {
entry.value_ptr.deinit(allocator);
}
self.pending.deinit(allocator);
}
fn request(ptr: *anyopaque, transfer: *Transfer) anyerror!void {
const self: *RobotsLayer = @ptrCast(@alignCast(ptr));
if (transfer.req.internal) {
return self.next.request(transfer);
}
const url = transfer.req.url;
const robots_url = try URL.getRobotsUrl(transfer.arena, url);
if (self.network.robot_store.get(robots_url)) |robot_entry| {
switch (robot_entry) {
.present => |robots| {
const path = URL.getPathname(url);
if (!robots.isAllowed(path)) {
try transfer.client.runNextTick(
transfer,
null,
.{
.run = struct {
fn run(t: *Transfer, _: ?*anyopaque) void {
defer t.deinit();
log.warn(.http, "blocked by robots", .{ .url = t.req.url });
t.req.error_callback(t.req.ctx, error.RobotsBlocked);
}
}.run,
},
);
return;
}
},
.absent => {},
}
return self.next.request(transfer);
}
return self.fetchRobotsThenRequest(robots_url, transfer);
}
fn fetchRobotsThenRequest(
self: *RobotsLayer,
robots_url: [:0]const u8,
transfer: *Transfer,
) !void {
const entry = try self.pending.getOrPut(self.allocator, robots_url);
if (!entry.found_existing) {
errdefer std.debug.assert(self.pending.remove(robots_url));
entry.value_ptr.* = .empty;
try entry.value_ptr.append(self.allocator, transfer);
transfer.park(.robots);
errdefer {
entry.value_ptr.deinit(self.allocator);
transfer.unpark();
}
const robots_ctx = try transfer.arena.create(RobotsContext);
robots_ctx.* = .{
.layer = self,
.buffer = .empty,
.arena = transfer.arena,
.robots_url = robots_url,
};
// CRITICAL: build a fresh Headers for the inner robots fetch.
// We value-copy req from the parent, but Headers is a struct wrapping
// a *curl_slist — value copy shares the pointer. Letting Client.request
// take ownership of a shared headers list means both transfers will
// free it at deinit time -> double-free. The robots.txt fetch is a
// system-level GET anyway, no need to inherit the parent's user headers.
var new_req = transfer.req;
new_req.headers = try transfer.client.newHeaders();
errdefer new_req.headers.deinit();
new_req.method = .GET;
new_req.url = robots_url;
new_req.internal = true;
new_req.resource_type = .fetch;
new_req.body = null;
new_req.ctx = robots_ctx;
new_req.start_callback = null;
new_req.header_callback = RobotsContext.headerCallback;
new_req.data_callback = RobotsContext.dataCallback;
new_req.done_callback = RobotsContext.doneCallback;
new_req.error_callback = RobotsContext.errorCallback;
new_req.shutdown_callback = RobotsContext.shutdownCallback;
log.debug(.browser, "fetching robots.txt", .{ .robots_url = robots_url });
try transfer.client.request(new_req, transfer.owner);
} else {
// Already one in flight, just queue behind.
try entry.value_ptr.append(self.allocator, transfer);
// Parked: RobotsLayer owns destruction via flushPending / flushPendingShutdown
// until robots.txt resolves. Without this, Client.request's errdefer (or
// any caller's cleanup) would deinit a transfer that's still on the
// pending list, leaving flushPending with a dangling pointer.
transfer.park(.robots);
}
}
fn flushPending(self: *RobotsLayer, robots_url: [:0]const u8, allowed: bool) void {
var queued = self.pending.fetchRemove(robots_url) orelse return;
defer queued.value.deinit(self.allocator);
for (queued.value.items) |transfer| {
if (!allowed) {
log.warn(.http, "blocked by robots", .{ .url = transfer.req.url });
transfer.abort(error.RobotsBlocked);
} else {
// Hand back to the layer chain. If a downstream layer commits
// (multi / queue / park), state advances past .created. If it
// fails before committing, we clean up here.
transfer.unpark();
self.next.request(transfer) catch |e| {
if (transfer.state == .created) {
transfer.abort(e);
}
};
}
}
}
// Invariant: shutdown_callback fires on a Transfer only via Transfer.kill,
// and the only callers of kill are Client.abortOwner / .abortRequests
// (owner-driven teardown). So if THIS robots fetch's shutdown_callback
// fired, the owner is being torn down — every parked transfer in this
// pending queue is on the same owner list and is already being killed by
// the same walk. We just need to drop the pending entry; the owner walk
// handles the rest. (If a future code path adds per-transfer kill
// without owner teardown, this assumption breaks — see comment above
// detachOrDeinit in HttpClient.zig.)
fn flushPendingShutdown(self: *RobotsLayer, robots_url: [:0]const u8) void {
var pending = self.pending.fetchRemove(robots_url) orelse return;
pending.value.deinit(self.allocator);
}
const RobotsContext = struct {
layer: *RobotsLayer,
arena: Allocator,
robots_url: [:0]const u8,
buffer: std.ArrayList(u8),
status: u16 = 0,
fn deinit(self: *RobotsContext) void {
self.buffer.deinit(self.layer.allocator);
self.layer.allocator.destroy(self);
}
fn headerCallback(response: Response) anyerror!HeaderResult {
const self: *RobotsContext = @ptrCast(@alignCast(response.ctx));
switch (response.inner) {
.transfer => |t| {
if (t.res.header) |hdr| {
log.debug(.browser, "robots status", .{ .status = hdr.status, .robots_url = self.robots_url });
self.status = hdr.status;
}
if (t.getContentLength()) |cl| {
try self.buffer.ensureTotalCapacity(self.arena, cl);
}
},
else => {},
}
return .proceed;
}
fn dataCallback(response: Response, data: []const u8) anyerror!void {
const self: *RobotsContext = @ptrCast(@alignCast(response.ctx));
try self.buffer.appendSlice(self.arena, data);
}
fn doneCallback(ctx_ptr: *anyopaque) anyerror!void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
const l = self.layer;
const robots_url = self.robots_url;
var allowed = true;
const network = l.network;
switch (self.status) {
200 => {
if (self.buffer.items.len > 0) {
const robots: ?Robots = network.robot_store.robotsFromBytes(
network.config.http_headers.user_agent,
self.buffer.items,
) catch blk: {
log.warn(.browser, "failed to parse robots", .{ .robots_url = robots_url });
try network.robot_store.putAbsent(robots_url);
break :blk null;
};
if (robots) |r| {
try network.robot_store.put(robots_url, r);
const path = URL.getPathname(l.pending.get(robots_url).?.items[0].req.url);
allowed = r.isAllowed(path);
}
}
},
404 => {
log.debug(.http, "robots not found", .{ .url = robots_url });
try network.robot_store.putAbsent(robots_url);
},
else => {
log.debug(.http, "unexpected status on robots", .{
.url = robots_url,
.status = self.status,
});
try network.robot_store.putAbsent(robots_url);
},
}
l.flushPending(robots_url, allowed);
}
fn errorCallback(ctx_ptr: *anyopaque, err: anyerror) void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
const l = self.layer;
const robots_url = self.robots_url;
log.warn(.http, "robots fetch failed", .{ .err = err });
l.flushPending(robots_url, true);
}
fn shutdownCallback(ctx_ptr: *anyopaque) void {
const self: *RobotsContext = @ptrCast(@alignCast(ctx_ptr));
const l = self.layer;
const robots_url = self.robots_url;
log.debug(.http, "robots fetch shutdown", .{});
l.flushPendingShutdown(robots_url);
}
};
-43
View File
@@ -1,43 +0,0 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const URL = @import("../../browser/URL.zig");
const Layer = @import("../../browser/HttpClient.zig").Layer;
const Transfer = @import("../../browser/HttpClient.zig").Transfer;
const WebBotAuthLayer = @This();
next: Layer = undefined,
pub fn layer(self: *WebBotAuthLayer) Layer {
return .{
.ptr = self,
.vtable = &.{ .request = request },
};
}
fn request(ptr: *anyopaque, transfer: *Transfer) anyerror!void {
const self: *WebBotAuthLayer = @ptrCast(@alignCast(ptr));
const wba = transfer.client.network.web_bot_auth orelse @panic("WebBotAuthLayer shouldn't be active without WebBotAuth");
const authority = URL.getHost(transfer.req.url);
try wba.signRequest(transfer.arena, &transfer.req.headers, authority);
return self.next.request(transfer);
}