commit e38628d9f89c1d136bf7e5d6393ea1af7efcad5a
parent b0f82196be5f87db02102408f8bb47d397648cb4
Author: triesap <tyson@radroots.org>
Date: Thu, 24 Sep 2026 01:48:30 +0000
H007: exact raw framing grammar and synthetic-error provenance (D43)
- parse the leading status line and exact line-delimited framing fields, rejecting malformed/duplicate/conflicting/overflow/unsupported declarations
- cap the header on bytes through the terminator and read exact declared bodies with explicit surplus/incomplete/overflow results
- require a real non-synthetic write failure to satisfy an expected peer close on both fixture serve paths
- add executed OB01-OB04 controls and a disclosed deterministic incremental read seam
Diffstat:
4 files changed, 1139 insertions(+), 124 deletions(-)
diff --git a/tests/bounded_call_helper.mojo b/tests/bounded_call_helper.mojo
@@ -76,6 +76,8 @@ def _field_allowed(key: String) -> Bool:
return True
if key == "body_bytes" or key == "body_match":
return True
+ if key == "surplus_bytes":
+ return True
return key == "declared_bytes" or key == "length_match"
@@ -84,7 +86,9 @@ def _field_numeric(key: String) -> Bool:
return True
if key == "head_ms" or key == "total_ms":
return True
- return key == "body_bytes" or key == "declared_bytes"
+ return (
+ key == "body_bytes" or key == "declared_bytes" or key == "surplus_bytes"
+ )
def _all_digits(text: String) -> Bool:
@@ -120,6 +124,7 @@ struct BoundedCallOutcome(Movable):
var body_match: String
var declared_bytes: Int
var length_match: String
+ var surplus_bytes: Int
def domain_failure(self) -> Bool:
return self.ok and self.outcome == "fail"
@@ -152,6 +157,8 @@ struct BoundedCallOutcome(Movable):
+ String(self.declared_bytes)
+ " length_match="
+ self.length_match
+ + " surplus_bytes="
+ + String(self.surplus_bytes)
)
@@ -172,6 +179,7 @@ def _empty_outcome(problem: String) -> BoundedCallOutcome:
body_match="",
declared_bytes=-1,
length_match="",
+ surplus_bytes=-1,
)
@@ -216,6 +224,7 @@ def parse_bounded_report(
var match_text = ""
var declared_text = ""
var length_text = ""
+ var surplus_text = ""
for index in range(len(keys)):
var key = keys[index]
var value = values[index]
@@ -245,6 +254,8 @@ def parse_bounded_report(
declared_text = value
elif key == "length_match":
length_text = value
+ elif key == "surplus_bytes":
+ surplus_text = value
if _field_numeric(key) and not _all_digits(value):
return _empty_outcome("report_non_numeric_" + key)
if kind == "" or correlation_text == "" or outcome == "":
@@ -284,6 +295,7 @@ def parse_bounded_report(
body_match=match_text,
declared_bytes=Int(declared_text) if declared_text != "" else -1,
length_match=length_text,
+ surplus_bytes=Int(surplus_text) if surplus_text != "" else -1,
)
@@ -337,24 +349,6 @@ def _raw_request_text(path: String) -> String:
)
-def _raw_status_code(head: String) raises -> Int:
- var prefix = "HTTP/1.1 "
- var start = head.find(prefix)
- if start < 0:
- return 0
- var base = start + prefix.byte_length()
- var length = 0
- var bytes = head.as_bytes()
- while base + length < len(bytes):
- var b = Int(bytes[base + length])
- if b < 48 or b > 57:
- break
- length += 1
- if length == 0:
- return 0
- return Int(String(head[byte = base : base + length]))
-
-
comptime RAW_MAX_HEADER_BYTES: Int = 65536
comptime RAW_MAX_BODY_BYTES: Int = 1048576
comptime RAW_READ_CHUNK_BYTES: Int = 1024
@@ -377,43 +371,170 @@ def _find_header_terminator(bytes: List[UInt8]) -> Int:
return -1
-def _declared_content_length(head_text: String) raises -> Int:
- """Lexical Content-Length from a decoded head, or -1 when absent.
+def _trim_ows_text(value: String) -> String:
+ """Trim only legal HTTP optional whitespace (SP / HTAB)."""
+ var start = 0
+ var end = value.byte_length()
+ var bytes = value.as_bytes()
+ while start < end and (Int(bytes[start]) == 32 or Int(bytes[start]) == 9):
+ start += 1
+ while end > start and (
+ Int(bytes[end - 1]) == 32 or Int(bytes[end - 1]) == 9
+ ):
+ end -= 1
+ if start == 0 and end == value.byte_length():
+ return String(value)
+ return String(value[byte=start:end])
+
+
+def _token_name(value: String) -> Bool:
+ """RFC 7230 token check for an exact header field name."""
+ if value.byte_length() == 0:
+ return False
+ for byte in value.as_bytes():
+ var b = Int(byte)
+ if b >= 48 and b <= 57:
+ continue
+ if b >= 65 and b <= 90:
+ continue
+ if b >= 97 and b <= 122:
+ continue
+ if (
+ b == 33
+ or b == 35
+ or b == 36
+ or b == 37
+ or b == 38
+ or b == 39
+ or b == 42
+ or b == 43
+ or b == 45
+ or b == 46
+ or b == 94
+ or b == 95
+ or b == 96
+ or b == 124
+ or b == 126
+ ):
+ continue
+ return False
+ return True
+
+
+def _value_legal(value: String) -> Bool:
+ """Reject control bytes in a header value (HTAB is the only legal one)."""
+ for byte in value.as_bytes():
+ var b = Int(byte)
+ if b == 9:
+ continue
+ if b < 32 or b == 127:
+ return False
+ return True
+
+
+def _ascii_digit(value: Int) -> Bool:
+ return value >= 48 and value <= 57
+
+
+@fieldwise_init
+struct HeadFraming(Movable):
+ """Exact framing parsed from one response head.
- Case-insensitive header name, ASCII digits only, with an explicit length
- guard. A malformed or oversized value is reported as absent rather than
- silently truncated; the caller treats that as an unknown declared size.
+ ``problem`` is an explicit bounded grammar failure token; otherwise
+ ``status`` is the leading three-digit status and ``declared`` is the exact
+ Content-Length (``-1`` when the header is absent).
"""
- var marker = "content-length"
- var m = marker.byte_length()
- var bytes = head_text.as_bytes()
- var n = len(bytes)
- var base = 0
- while base + m <= n:
- var matched = True
- for k in range(m):
- var hb = Int(bytes[base + k])
- if hb >= 65 and hb <= 90:
- hb += 32
- if hb != Int(marker.as_bytes()[k]):
- matched = False
- break
- if matched:
- var i = base + m
- while i < n and (Int(bytes[i]) == 32 or Int(bytes[i]) == 9):
- i += 1
- if i < n and Int(bytes[i]) == 58:
- i += 1
- while i < n and (Int(bytes[i]) == 32 or Int(bytes[i]) == 9):
- i += 1
- var j = i
- while j < n and Int(bytes[j]) >= 48 and Int(bytes[j]) <= 57:
- j += 1
- if j == i or j - i > 9:
- return -1
- return Int(String(head_text[byte=i:j]))
- base += 1
- return -1
+
+ var problem: String
+ var status: Int
+ var declared: Int
+
+
+def _parse_response_head(head_text: String) raises -> HeadFraming:
+ """Parse the leading HTTP/1.1 status line and exact framing headers.
+
+ OB01: the raw observer accepts only an exact leading
+ ``HTTP/1.1|HTTP/1.0 SP three-digit-status`` line and line-delimited
+ ``name: value`` headers. A malformed, duplicated, conflicting, negative,
+ junk or overflowing framing declaration and an unsupported
+ Transfer-Encoding are explicit bounded problem tokens. A status-looking
+ string inside a header value is never read as the response status, and a
+ different field name (``X-Content-Length``) is preserved as an unknown
+ header rather than interpreted as ``Content-Length``.
+ """
+ var lines = head_text.split("\r\n")
+ if len(lines) < 1:
+ return HeadFraming("raw_status_missing", 0, -1)
+ var status_line = String(lines[0])
+ var version = ""
+ if status_line.startswith("HTTP/1.1 "):
+ version = "HTTP/1.1 "
+ elif status_line.startswith("HTTP/1.0 "):
+ version = "HTTP/1.0 "
+ else:
+ return HeadFraming("raw_status_missing", 0, -1)
+ var rest = String(status_line[byte = version.byte_length() :])
+ if rest.byte_length() < 3:
+ return HeadFraming("raw_status_malformed", 0, -1)
+ var digits = String(rest[byte=0:3])
+ for byte in digits.as_bytes():
+ if not _ascii_digit(Int(byte)):
+ return HeadFraming("raw_status_malformed", 0, -1)
+ if rest.byte_length() > 3:
+ if Int(rest.as_bytes()[3]) != 32:
+ return HeadFraming("raw_status_malformed", 0, -1)
+ var status_value = Int(digits)
+ if status_value < 100 or status_value > 599:
+ # A three-digit but out-of-range status is malformed, never a
+ # successfully observed call.
+ return HeadFraming("raw_status_malformed", 0, -1)
+ var declared = -1
+ var have_length = False
+ var have_transfer = False
+ for index in range(1, len(lines)):
+ var line = String(lines[index])
+ if line.byte_length() == 0:
+ continue
+ var first = Int(line.as_bytes()[0])
+ if first == 32 or first == 9:
+ # obs-fold / a continuation line is not a legal standalone header.
+ return HeadFraming("raw_header_malformed", 0, -1)
+ var colon = line.find(":")
+ if colon <= 0:
+ return HeadFraming("raw_header_malformed", 0, -1)
+ var name = String(line[byte=0:colon])
+ if not _token_name(name):
+ return HeadFraming("raw_header_malformed", 0, -1)
+ var value = _trim_ows_text(String(line[byte = colon + 1 :]))
+ if not _value_legal(value):
+ return HeadFraming("raw_header_malformed", 0, -1)
+ var lower = name.lower()
+ if lower == "content-length":
+ if have_length:
+ # Includes an identical duplicate: a second framing field is
+ # ambiguous even when the value agrees.
+ return HeadFraming("raw_content_length_duplicate", 0, -1)
+ if value.byte_length() == 0:
+ return HeadFraming("raw_content_length_malformed", 0, -1)
+ for byte in value.as_bytes():
+ if not _ascii_digit(Int(byte)):
+ return HeadFraming("raw_content_length_malformed", 0, -1)
+ if value.byte_length() > 9:
+ return HeadFraming("raw_content_length_overflow", 0, -1)
+ var parsed = Int(value)
+ if parsed > RAW_MAX_BODY_BYTES:
+ return HeadFraming("raw_body_overflow", 0, -1)
+ declared = parsed
+ have_length = True
+ elif lower == "transfer-encoding":
+ if have_transfer:
+ return HeadFraming("raw_transfer_encoding_duplicate", 0, -1)
+ have_transfer = True
+ if value.lower() != "identity":
+ return HeadFraming("raw_transfer_encoding_unsupported", 0, -1)
+ if have_transfer and have_length:
+ return HeadFraming("raw_framing_conflict", 0, -1)
+ return HeadFraming("", status_value, declared)
def _bytes_contain_from(bytes: List[UInt8], start: Int, needle: String) -> Bool:
@@ -456,6 +577,7 @@ struct RawAccounting(Movable):
var body_match: String
var declared_bytes: Int
var length_match: String
+ var surplus_bytes: Int
def account_raw_bytes(raw: List[UInt8], expect: String) raises -> RawAccounting:
@@ -463,12 +585,16 @@ def account_raw_bytes(raw: List[UInt8], expect: String) raises -> RawAccounting:
The buffer is treated as one byte stream: packet/read boundaries are never
protocol boundaries, and the body begins exactly after CRLFCRLF. A valid
- declared Content-Length is compared with the observed body byte count.
+ declared Content-Length is compared with the observed body byte count and
+ any already-buffered bytes beyond it are reported explicitly as surplus.
+ Malformed, duplicated, conflicting, negative, junk or overflowing framing
+ and an unsupported Transfer-Encoding are explicit problem tokens rather
+ than silently absent fields.
"""
var head_end = _find_header_terminator(raw)
if head_end < 0:
return RawAccounting(
- "raw_head_incomplete", 0, 0, "unknown", -1, "unknown"
+ "raw_head_incomplete", 0, 0, "unknown", -1, "unknown", 0
)
var head_bytes = List[UInt8]()
for index in range(head_end):
@@ -482,41 +608,59 @@ def account_raw_bytes(raw: List[UInt8], expect: String) raises -> RawAccounting:
except:
decoded = False
if not decoded:
- return RawAccounting("raw_head_decode", 0, 0, "unknown", -1, "unknown")
+ return RawAccounting(
+ "raw_head_decode", 0, 0, "unknown", -1, "unknown", 0
+ )
+ var framing = _parse_response_head(head_text)
+ if framing.problem != "":
+ return RawAccounting(framing.problem, 0, 0, "unknown", -1, "unknown", 0)
var body_bytes = len(raw) - head_end
- var declared = _declared_content_length(head_text)
var match_text = "unknown"
if expect != "":
match_text = "yes" if _bytes_contain_from(
raw, head_end, expect
) else "no"
var length_match = "unknown"
- if declared >= 0:
- length_match = "yes" if body_bytes == declared else "no"
- var status = _raw_status_code(head_text)
- if status == 0:
- # A header block with no parseable HTTP/1.1 status line is a malformed
- # observation, never a successful call with an invented status.
- return RawAccounting(
- "raw_status_missing",
- 0,
- body_bytes,
- match_text,
- declared,
- length_match,
- )
+ var surplus = 0
+ if framing.declared >= 0:
+ length_match = "yes" if body_bytes == framing.declared else "no"
+ if body_bytes > framing.declared:
+ surplus = body_bytes - framing.declared
return RawAccounting(
"",
- status,
+ framing.status,
body_bytes,
match_text,
- declared,
+ framing.declared,
length_match,
+ surplus,
)
+def _plan_read_size(plan: List[Int], index: Int, default: Int) -> Int:
+ """Deterministic bounded read size from an optional chunk plan (OB02).
+
+ The plan element at ``index`` caps one actual read; once the plan is
+ exhausted the last element stays in force. An absent plan uses ``default``.
+ This is a disclosed synthetic acquisition seam, used only to force a header
+ terminator or a multi-byte body character to split across real incremental
+ reads; it is not a claim about packet boundaries.
+ """
+ if len(plan) == 0:
+ return default
+ var i = index if index < len(plan) else len(plan) - 1
+ var want = plan[i]
+ if want <= 0:
+ want = 1
+ return want if want < default else default
+
+
def _child_raw_head_body(
- head: String, port: Int, path: String, expect: String
+ head: String,
+ port: Int,
+ path: String,
+ expect: String,
+ chunk_plan: List[Int],
) raises -> String:
"""Owned raw client: record head/body arrival timing for a scripted peer.
@@ -524,96 +668,142 @@ def _child_raw_head_body(
parsed response, so a headers-before-body-stall claim can be proved from the
wire while still running under the parent's finite deadline.
- RP01: bytes coalesced with the header terminator stay with the body, the
- accumulation is byte-oriented and bounded, and decode happens once over the
- complete buffer so a split multi-byte character is never misread as a
- packet boundary. A valid declared Content-Length bounds the read; the
- observed body byte count is compared with it and reported truthfully.
+ OB01/OB02: the head grammar is parsed exactly, the header cap counts bytes
+ through the terminator only (a coalesced body is not charged to it), a valid
+ declared Content-Length is read as an exact message body with any surplus
+ reported explicitly, a short read to EOF is an explicit incomplete
+ observation, and a no-length response completes only at EOF inside the
+ body cap. ``chunk_plan`` is a disclosed deterministic read seam that runs
+ the same incremental loop.
"""
var client = TcpStream.connect(SocketAddr.localhost(UInt16(port)))
var start = now_ms()
client.write_all(Span[UInt8, _](_raw_request_text(path).as_bytes()))
var raw = List[UInt8]()
var buffer = InlineArray[Byte, RAW_READ_CHUNK_BYTES](fill=0)
+ var plan_index = 0
var head_end = -1
while head_end < 0:
- var n = client.read(buffer.unsafe_ptr(), RAW_READ_CHUNK_BYTES)
+ var want = _plan_read_size(chunk_plan, plan_index, RAW_READ_CHUNK_BYTES)
+ plan_index += 1
+ var n = client.read(buffer.unsafe_ptr(), want)
if n <= 0:
client.close()
return head + "outcome=fail cause=raw_eof reason=head_eof"
for index in range(n):
raw.append(UInt8(Int(buffer[index])))
- if len(raw) > RAW_MAX_HEADER_BYTES:
+ head_end = _find_header_terminator(raw)
+ if head_end < 0 and len(raw) > RAW_MAX_HEADER_BYTES:
client.close()
return (
head
+ "outcome=fail cause=raw_head_overflow reason=head_overflow"
)
- head_end = _find_header_terminator(raw)
+ if head_end > RAW_MAX_HEADER_BYTES:
+ client.close()
+ return (
+ head + "outcome=fail cause=raw_head_overflow reason=head_overflow"
+ )
var head_ms = now_ms() - start
- var declared = -1
var head_bytes = List[UInt8]()
for index in range(head_end):
head_bytes.append(raw[index])
+ var head_text = ""
var decoded = True
try:
- declared = _declared_content_length(
- String(
- from_utf8=Span(
- ptr=head_bytes.unsafe_ptr(), length=len(head_bytes)
- )
- )
+ head_text = String(
+ from_utf8=Span(ptr=head_bytes.unsafe_ptr(), length=len(head_bytes))
)
except:
decoded = False
if not decoded:
client.close()
return head + "outcome=fail cause=raw_head_decode reason=head_decode"
+ var framing = _parse_response_head(head_text)
+ if framing.problem != "":
+ client.close()
+ return (
+ head
+ + "outcome=fail cause="
+ + framing.problem
+ + " reason="
+ + framing.problem
+ )
+ var declared = framing.declared
+ var status = framing.status
var body_bytes = len(raw) - head_end
+ var incomplete = False
+ var overflow = False
while True:
if declared >= 0 and body_bytes >= declared:
break
- if body_bytes >= RAW_MAX_BODY_BYTES:
+ if declared < 0 and body_bytes >= RAW_MAX_BODY_BYTES:
+ overflow = True
+ break
+ var want2 = _plan_read_size(
+ chunk_plan, plan_index, RAW_READ_CHUNK_BYTES
+ )
+ plan_index += 1
+ var remaining = RAW_MAX_BODY_BYTES - body_bytes
+ if want2 > remaining:
+ want2 = remaining
+ if want2 <= 0:
+ overflow = True
break
- var want = min(RAW_READ_CHUNK_BYTES, RAW_MAX_BODY_BYTES - body_bytes)
- var n2 = client.read(buffer.unsafe_ptr(), want)
+ var n2 = client.read(buffer.unsafe_ptr(), want2)
if n2 <= 0:
+ # A length-delimited response that ends before its declared body is
+ # explicitly incomplete; a no-length response completes at EOF.
+ if declared >= 0 and body_bytes < declared:
+ incomplete = True
break
for index in range(n2):
raw.append(UInt8(Int(buffer[index])))
body_bytes += n2
var total_ms = now_ms() - start
client.close()
- var accounting = account_raw_bytes(raw^, expect)
- if accounting.problem != "":
- return (
- head
- + "outcome=fail cause="
- + accounting.problem
- + " reason="
- + accounting.problem
- )
+ var match_text = "unknown"
+ if expect != "":
+ match_text = "yes" if _bytes_contain_from(
+ raw, head_end, expect
+ ) else "no"
+ var length_match = "unknown"
+ var surplus = 0
+ if declared >= 0:
+ length_match = "yes" if body_bytes == declared else "no"
+ if body_bytes > declared:
+ surplus = body_bytes - declared
var declared_field = ""
- if accounting.declared_bytes >= 0:
- # A negative declared length is omitted rather than emitted, because the
- # report grammar accepts only non-negative numeric fields.
- declared_field = " declared_bytes=" + String(accounting.declared_bytes)
- return (
- head
- + "outcome=ok status="
- + String(accounting.status)
- + " head_ms="
+ if declared >= 0:
+ declared_field = " declared_bytes=" + String(declared)
+ var tail = (
+ " head_ms="
+ String(head_ms)
+ " total_ms="
+ String(total_ms)
+ " body_bytes="
- + String(accounting.body_bytes)
+ + String(body_bytes)
+ " body_match="
- + accounting.body_match
+ + match_text
+ " length_match="
- + accounting.length_match
+ + length_match
+ + " surplus_bytes="
+ + String(surplus)
+ declared_field
)
+ if incomplete:
+ return (
+ head
+ + "outcome=fail cause=raw_body_incomplete reason=body_eof"
+ + tail
+ )
+ if overflow:
+ return (
+ head
+ + "outcome=fail cause=raw_body_overflow reason=body_overflow"
+ + tail
+ )
+ return head + "outcome=ok status=" + String(status) + tail
def _child_report(
@@ -623,6 +813,7 @@ def _child_report(
correlation: Int,
raw_path: String,
raw_expect: String,
+ raw_chunk_plan: List[Int],
) raises -> String:
"""One bounded report line from inside the forked provider-call child."""
var head = _report_prefix(kind, correlation)
@@ -666,7 +857,9 @@ def _child_report(
)
return head + "outcome=ok status=" + String(response.status)
if kind == "raw_head_body" or kind == "raw_jev_head_body":
- return _child_raw_head_body(head, port, raw_path, raw_expect)
+ return _child_raw_head_body(
+ head, port, raw_path, raw_expect, raw_chunk_plan
+ )
return head + "outcome=fail cause=unknown_kind reason=unknown_kind"
@@ -699,6 +892,7 @@ struct BoundedCallReport(Movable):
var body_match: String
var declared_bytes: Int
var length_match: String
+ var surplus_bytes: Int
var problem: String
var report: String
var elapsed_ms: Int
@@ -739,6 +933,8 @@ struct BoundedCallReport(Movable):
+ String(self.declared_bytes)
+ " length_match="
+ self.length_match
+ + " surplus_bytes="
+ + String(self.surplus_bytes)
+ " problem="
+ (self.problem if self.problem != "" else "-")
+ " elapsed_ms="
@@ -765,6 +961,7 @@ def run_bounded_call(
fault_wait_errors: Int = 0,
fault_poll_eintrs: Int = 0,
fault_poll_errors: Int = 0,
+ raw_chunk_plan: List[Int] = List[Int](),
) raises -> BoundedCallReport:
"""Run one risky provider call under a parent-enforced finite deadline.
@@ -859,7 +1056,13 @@ def run_bounded_call(
else:
try:
payload = _child_report(
- kind, port, timeout_ms, correlation, raw_path, raw_expect
+ kind,
+ port,
+ timeout_ms,
+ correlation,
+ raw_path,
+ raw_expect,
+ raw_chunk_plan,
)
except e:
payload = (
@@ -983,6 +1186,7 @@ def run_bounded_call(
body_match=parsed.body_match,
declared_bytes=parsed.declared_bytes,
length_match=parsed.length_match,
+ surplus_bytes=parsed.surplus_bytes,
problem=problem,
report=String(report_text.strip()),
elapsed_ms=elapsed_ms,
diff --git a/tests/strict_fixture.mojo b/tests/strict_fixture.mojo
@@ -909,8 +909,15 @@ def serve_scripts(
) if synthetic else "realerrno"
+ String(raw_errno)
)
+ # OB03: provenance is kept separate from the error class. A
+ # synthetic/injected seam is never a real peer close, so it
+ # cannot satisfy an expected real close even when it maps to
+ # the same bounded class (EPIPE -> broken_pipe, ECONNRESET ->
+ # peer_reset). It is rejected here with its synthetic evidence
+ # label instead of being accepted.
if (
script.expect_peer_close
+ and not synthetic
and is_peer_close_cause(script.expected_close_cause)
and observed == script.expected_close_cause
and phase == script.expected_close_phase
diff --git a/tests/test_jev.mojo b/tests/test_jev.mojo
@@ -522,6 +522,11 @@ comptime JEV_CORRELATION_RAW_TRUNCATED = 204
# H007 RP03: distinct correlations for the Jev write-error controls.
comptime JEV_CORRELATION_SYNTHETIC_DESCRIPTOR = 205
comptime JEV_CORRELATION_REAL_ERRNO = 206
+# H007 OB01-OB03: period-14 observation-integrity controls on the Jev path.
+comptime JEV_CORRELATION_OB_SPLIT_TERMINATOR = 210
+comptime JEV_CORRELATION_OB_SPLIT_BODY = 211
+comptime JEV_CORRELATION_OB_SURPLUS = 212
+comptime JEV_CORRELATION_OB_NO_LENGTH = 213
def _raw_jev_request_text(path: String) -> String:
@@ -670,12 +675,15 @@ def test_jev_raw_head_body_reports_truncated_length_mismatch() raises:
"/v1/systemone",
"SHORT",
)
- assert_true(report.ok())
- assert_equal(report.status, 200)
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_body_incomplete")
+ assert_equal(report.status, 0)
assert_equal(report.body_bytes, 5)
assert_equal(report.body_match, "yes")
assert_equal(report.declared_bytes, 20)
assert_equal(report.length_match, "no")
+ assert_equal(report.surplus_bytes, 0)
started.stub.wait()
assert_true(started.stub.ok())
guard.assert_clean()
@@ -1092,3 +1100,240 @@ def test_jev_wire_attempt_counts_for_non_retryable_failure() raises:
assert_equal(started.stub.request_count(), 1)
assert_equal(started.stub.connection_count(), 1)
guard.assert_clean()
+
+
+# OB02: the Jev raw path must execute the same exact framing/cap/surplus and
+# actual incremental fragmentation controls, not only a dormant shared branch.
+
+
+def test_jev_raw_head_body_split_header_terminator_is_preserved() raises:
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "jev_split_terminator", "POST", "/v1/systemone", 200, "JEVBODY"
+ )
+ )
+ var plan = List[Int]()
+ plan.append(1)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_OB_SPLIT_TERMINATOR,
+ "/v1/systemone",
+ "JEVBODY",
+ raw_chunk_plan=plan,
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 7)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 7)
+ assert_equal(report.length_match, "yes")
+ assert_equal(report.surplus_bytes, 0)
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ guard.assert_clean()
+
+
+def test_jev_raw_head_body_split_multibyte_body_is_preserved() raises:
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "jev_split_multibyte",
+ "POST",
+ "/v1/systemone",
+ 200,
+ '{"mark":"x☃y"}',
+ )
+ )
+ var plan = List[Int]()
+ plan.append(1)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_OB_SPLIT_BODY,
+ "/v1/systemone",
+ "x☃y",
+ raw_chunk_plan=plan,
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, '{"mark":"x☃y"}'.byte_length())
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.length_match, "yes")
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ guard.assert_clean()
+
+
+def test_jev_raw_head_body_accounts_buffered_surplus() raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_buffered_surplus", "POST", "/v1/systemone", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-length: 3\r\nconnection: close\r\n\r\nabcde"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_OB_SURPLUS,
+ "/v1/systemone",
+ "abc",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 5)
+ assert_equal(report.declared_bytes, 3)
+ assert_equal(report.length_match, "no")
+ assert_equal(report.surplus_bytes, 2)
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ guard.assert_clean()
+
+
+def test_jev_raw_head_body_no_length_completes_at_eof() raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_no_length", "POST", "/v1/systemone", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\nconnection: close\r\n\r\nNOLENGTHBODY"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ var report = run_bounded_call(
+ "raw_jev_head_body",
+ started.port,
+ 5000,
+ 5000,
+ guard,
+ JEV_CORRELATION_OB_NO_LENGTH,
+ "/v1/systemone",
+ "NOLENGTHBODY",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, "NOLENGTHBODY".byte_length())
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, -1)
+ assert_equal(report.length_match, "unknown")
+ started.stub.wait()
+ assert_true(started.stub.ok())
+ guard.assert_clean()
+
+
+# OB03: the Jev serve path must also reject every synthetic write-error seam,
+# with and without an expected-close declaration.
+
+
+def _assert_jev_synthetic_errno_rejected(errno: Int, declared: Bool) raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_synthetic_errno", "POST", "/v1/systemone", 200, '{"ok":true}'
+ )
+ script.stall_after_head_ms = 200
+ if declared:
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "body_stall"
+ script.inject_write_errno = errno
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ _ = _raw_jev_send_and_read(started.port, "/v1/systemone")
+ started.stub.reap()
+ assert_true(not started.stub.ok())
+ assert_equal(started.stub.phase(), "peer_close")
+ assert_true(started.stub.reason().find("unexpected_write_") >= 0)
+ assert_true(started.stub.reason().find("synthdecl") >= 0)
+ guard.assert_clean()
+
+
+def test_jev_scripted_epipe_with_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EPIPE.value), True)
+
+
+def test_jev_scripted_epipe_without_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EPIPE.value), False)
+
+
+def test_jev_scripted_econnreset_with_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.ECONNRESET.value), True)
+
+
+def test_jev_scripted_econnreset_without_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.ECONNRESET.value), False)
+
+
+def test_jev_scripted_eagain_with_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EAGAIN.value), True)
+
+
+def test_jev_scripted_eagain_without_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EAGAIN.value), False)
+
+
+def test_jev_scripted_ebadf_with_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EBADF.value), True)
+
+
+def test_jev_scripted_ebadf_without_declaration_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_rejected(Int(ErrNo.EBADF.value), False)
+
+
+def _assert_jev_synthetic_errno_delayed_write_rejected(errno: Int) raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "jev_synthetic_delayed", "POST", "/v1/systemone", 200, '{"ok":true}'
+ )
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "delayed_write"
+ script.inject_write_errno = errno
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_jev_scripted_auto(scripts^, guard) as started:
+ _ = _raw_jev_send_and_read(started.port, "/v1/systemone")
+ started.stub.reap()
+ assert_true(not started.stub.ok())
+ assert_equal(started.stub.phase(), "peer_close")
+ assert_true(started.stub.reason().find("unexpected_write_") >= 0)
+ assert_true(started.stub.reason().find("_delayed_write_") >= 0)
+ assert_true(started.stub.reason().find("synthdecl") >= 0)
+ guard.assert_clean()
+
+
+def test_jev_scripted_epipe_delayed_write_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_delayed_write_rejected(Int(ErrNo.EPIPE.value))
+
+
+def test_jev_scripted_econnreset_delayed_write_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_delayed_write_rejected(
+ Int(ErrNo.ECONNRESET.value)
+ )
+
+
+def test_jev_scripted_eagain_delayed_write_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_delayed_write_rejected(Int(ErrNo.EAGAIN.value))
+
+
+def test_jev_scripted_ebadf_delayed_write_is_not_peer_close() raises:
+ _assert_jev_synthetic_errno_delayed_write_rejected(Int(ErrNo.EBADF.value))
diff --git a/tests/test_provider_adapter.mojo b/tests/test_provider_adapter.mojo
@@ -86,6 +86,18 @@ comptime BOUNDED_CORRELATION_DECLARED_VS_SYNTHETIC = 131
# H007 RP01: malformed/error raw-observation controls.
comptime BOUNDED_CORRELATION_RAW_EOF = 132
comptime BOUNDED_CORRELATION_RAW_MALFORMED = 133
+# H007 OB01-OB03: period-14 observation-integrity controls.
+comptime BOUNDED_CORRELATION_OB_UNKNOWN_LENGTH = 140
+comptime BOUNDED_CORRELATION_OB_JUNK_LENGTH = 141
+comptime BOUNDED_CORRELATION_OB_DUP_LENGTH = 142
+comptime BOUNDED_CORRELATION_OB_OVER_CAP = 143
+comptime BOUNDED_CORRELATION_OB_SURPLUS = 144
+comptime BOUNDED_CORRELATION_OB_NO_LENGTH = 145
+comptime BOUNDED_CORRELATION_OB_SPLIT_TERMINATOR = 146
+comptime BOUNDED_CORRELATION_OB_SPLIT_BODY = 147
+comptime BOUNDED_CORRELATION_OB_NEAR_CAP = 148
+comptime BOUNDED_CORRELATION_OB_MALFORMED_HEAD = 149
+comptime BOUNDED_CORRELATION_OB_SYNTHETIC = 150
from flare.net import SocketAddr
from flare.tcp import TcpStream
@@ -712,10 +724,10 @@ def test_provider_raw_head_body_accounts_empty_body() raises:
def test_provider_raw_head_body_reports_truncated_length_mismatch() raises:
- # RP01: an explicit truncation control. The peer declares twenty body bytes
- # but sends five and closes; the observer must report the exact five
- # observed bytes, the twenty declared and a length mismatch, never the
- # declared length as if it had arrived.
+ # OB02: an explicit truncation control. The peer declares twenty body bytes
+ # but sends five and closes; the observer must report an explicit
+ # incomplete observation with the exact five observed bytes and the twenty
+ # declared, never an unqualified successful call.
var scripts = List[ExchangeScript]()
var script = exchange_script(
"truncated_body", "POST", "/v1/chat/completions", 200, ""
@@ -737,12 +749,15 @@ def test_provider_raw_head_body_reports_truncated_length_mismatch() raises:
"/v1/chat/completions",
"SHORT",
)
- assert_true(report.ok())
- assert_equal(report.status, 200)
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_body_incomplete")
+ assert_equal(report.status, 0)
assert_equal(report.body_bytes, 5)
assert_equal(report.body_match, "yes")
assert_equal(report.declared_bytes, 20)
assert_equal(report.length_match, "no")
+ assert_equal(report.surplus_bytes, 0)
provider_stub.wait()
assert_true(provider_stub.ok())
guard.assert_clean()
@@ -1843,5 +1858,549 @@ def test_query_analysis_boundary_characterizes_duplicate_and_null_fields() raise
assert_true(terms_message.find("provider_schema_invalid") >= 0)
+def _raw_bytes(text: String) -> List[UInt8]:
+ var raw = List[UInt8]()
+ for byte in text.as_bytes():
+ raw.append(UInt8(Int(byte)))
+ return raw^
+
+
+def _repeat_text(mark: String, count: Int) -> String:
+ var out = List[UInt8]()
+ var mark_bytes = mark.as_bytes()
+ for _ in range(count):
+ for index in range(len(mark_bytes)):
+ out.append(UInt8(Int(mark_bytes[index])))
+ return String(unsafe_from_utf8=Span(ptr=out.unsafe_ptr(), length=len(out)))
+
+
+# OB01: the exact raw response grammar. A status-looking value inside a header
+# is never the status; a different field name is preserved as unknown; junk,
+# duplicate, conflicting, negative, overflowing or unsupported framing is an
+# explicit bounded problem rather than a silently absent field.
+
+
+def test_raw_accounting_preserves_unknown_length_like_field() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\nX-Content-Length: 3\r\n\r\nabc"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "")
+ assert_equal(accounting.status, 200)
+ assert_equal(accounting.body_bytes, 3)
+ assert_equal(accounting.declared_bytes, -1)
+ assert_equal(accounting.length_match, "unknown")
+ assert_equal(accounting.surplus_bytes, 0)
+
+
+def test_raw_accounting_rejects_junk_content_length() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\nContent-Length: 3junk\r\n\r\nabc"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_content_length_malformed")
+
+
+def test_raw_accounting_rejects_identical_duplicate_length() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes(
+ "HTTP/1.1 200 OK\r\nContent-Length: 3\r\n"
+ "Content-Length: 3\r\n\r\nabc"
+ ),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_content_length_duplicate")
+
+
+def test_raw_accounting_rejects_conflicting_duplicate_length() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes(
+ "HTTP/1.1 200 OK\r\nContent-Length: 3\r\n"
+ "Content-Length: 9\r\n\r\nabc"
+ ),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_content_length_duplicate")
+
+
+def test_raw_accounting_rejects_garbage_status_with_header_status_value() raises:
+ # A status-looking string in a header after a garbage first line must never
+ # be read as the response status.
+ var accounting = account_raw_bytes(
+ _raw_bytes(
+ "garbage\r\nx-note: HTTP/1.1 200 OK\r\nContent-Length: 3\r\n\r\nabc"
+ ),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_status_missing")
+ assert_equal(accounting.status, 0)
+
+
+def test_raw_accounting_rejects_unsupported_transfer_encoding() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\ntransfer-encoding: chunked\r\n\r\nabc"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_transfer_encoding_unsupported")
+
+
+def test_raw_accounting_rejects_negative_content_length() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\nContent-Length: -3\r\n\r\nabc"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_content_length_malformed")
+
+
+def test_raw_accounting_rejects_overflow_content_length() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes(
+ "HTTP/1.1 200 OK\r\nContent-Length: 99999999999999999999\r\n\r\nabc"
+ ),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_content_length_overflow")
+
+
+def test_raw_accounting_rejects_malformed_header_line() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\nbad header line\r\n\r\nabc"), "abc"
+ )
+ assert_equal(accounting.problem, "raw_header_malformed")
+
+
+def test_raw_accounting_rejects_out_of_range_status() raises:
+ # OB01: only a valid three-digit HTTP status (100..599) is accepted.
+ var low = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 000 X\r\ncontent-length: 3\r\n\r\nabc"), "abc"
+ )
+ assert_equal(low.problem, "raw_status_malformed")
+ var high = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 999 X\r\ncontent-length: 3\r\n\r\nabc"), "abc"
+ )
+ assert_equal(high.problem, "raw_status_malformed")
+
+
+def test_raw_accounting_accepts_identity_transfer_encoding() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\ntransfer-encoding: identity\r\n\r\nabc"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "")
+ assert_equal(accounting.status, 200)
+ assert_equal(accounting.declared_bytes, -1)
+ assert_equal(accounting.length_match, "unknown")
+
+
+def test_raw_accounting_rejects_transfer_with_length_conflict() raises:
+ var accounting = account_raw_bytes(
+ _raw_bytes(
+ "HTTP/1.1 200 OK\r\ntransfer-encoding: identity\r\n"
+ "content-length: 3\r\n\r\nabc"
+ ),
+ "abc",
+ )
+ assert_equal(accounting.problem, "raw_framing_conflict")
+
+
+def test_raw_accounting_reports_buffered_surplus() raises:
+ # OB02: bytes already buffered beyond the declared body are reported
+ # explicitly instead of being silently folded into the declared length.
+ var accounting = account_raw_bytes(
+ _raw_bytes("HTTP/1.1 200 OK\r\ncontent-length: 3\r\n\r\nabcde"),
+ "abc",
+ )
+ assert_equal(accounting.problem, "")
+ assert_equal(accounting.body_bytes, 5)
+ assert_equal(accounting.declared_bytes, 3)
+ assert_equal(accounting.length_match, "no")
+ assert_equal(accounting.surplus_bytes, 2)
+
+
+# OB02 acquisition controls through the actual incremental read path.
+
+
+def test_provider_raw_head_body_accounts_buffered_surplus() raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "buffered_surplus", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-length: 3\r\nconnection: close\r\n\r\nabcde"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_SURPLUS,
+ "/v1/chat/completions",
+ "abc",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 5)
+ assert_equal(report.declared_bytes, 3)
+ assert_equal(report.length_match, "no")
+ assert_equal(report.surplus_bytes, 2)
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_rejects_declared_over_cap() raises:
+ # OB02: a declared length beyond the body cap is an explicit overflow, never
+ # an accepted observation with a silently truncated prefix.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "declared_over_cap", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-length: 99999999\r\n"
+ "connection: close\r\n\r\nabc"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_OVER_CAP,
+ "/v1/chat/completions",
+ "abc",
+ )
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_body_overflow")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_rejects_duplicate_length() raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "duplicate_length", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-length: 3\r\ncontent-length: 3\r\n"
+ "connection: close\r\n\r\nabc"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_DUP_LENGTH,
+ "/v1/chat/completions",
+ "abc",
+ )
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_content_length_duplicate")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_no_length_completes_at_eof() raises:
+ # OB02: a response without a declared length completes at EOF inside the
+ # cap, with an explicitly unknown length rather than an invented one.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "no_length", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\nconnection: close\r\n\r\nNOLENGTHBODY"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_NO_LENGTH,
+ "/v1/chat/completions",
+ "NOLENGTHBODY",
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, "NOLENGTHBODY".byte_length())
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, -1)
+ assert_equal(report.length_match, "unknown")
+ assert_equal(report.surplus_bytes, 0)
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_split_header_terminator_is_preserved() raises:
+ # OB02: the real incremental acquisition path is driven with one-byte reads,
+ # so the CRLFCRLF terminator is split across reads. Concatenating loops
+ # before a pure parser call is not fragmentation proof, so the seam caps the
+ # actual read.
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "split_terminator", "POST", "/v1/chat/completions", 200, "BODYMARK"
+ )
+ )
+ var plan = List[Int]()
+ plan.append(1)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_SPLIT_TERMINATOR,
+ "/v1/chat/completions",
+ "BODYMARK",
+ raw_chunk_plan=plan,
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 8)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 8)
+ assert_equal(report.length_match, "yes")
+ assert_equal(report.surplus_bytes, 0)
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_split_multibyte_body_is_preserved() raises:
+ # OB02: a one-byte read seam splits the three-byte UTF-8 character across
+ # reads; the byte accumulation must still preserve it whole.
+ var scripts = List[ExchangeScript]()
+ scripts.append(
+ exchange_script(
+ "split_multibyte",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"mark":"a☃b"}',
+ )
+ )
+ var plan = List[Int]()
+ plan.append(1)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_SPLIT_BODY,
+ "/v1/chat/completions",
+ "a☃b",
+ raw_chunk_plan=plan,
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, '{"mark":"a☃b"}'.byte_length())
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.length_match, "yes")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_near_cap_coalesced_body_is_not_head_overflow() raises:
+ # OB02: the header cap counts bytes through the terminator only. A valid
+ # near-cap header whose final 1000-byte read also carries coalesced body
+ # bytes must not be reported as a header overflow.
+ var pad = _repeat_text("a", 65459)
+ var body = "ENDMARK" + _repeat_text("y", 593)
+ var raw = (
+ "HTTP/1.1 200 OK\r\nx-pad: "
+ + pad
+ + "\r\ncontent-length: 600\r\nconnection: close\r\n\r\n"
+ + body
+ )
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "near_cap", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = raw
+ scripts.append(script^)
+ var plan = List[Int]()
+ plan.append(1000)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_NEAR_CAP,
+ "/v1/chat/completions",
+ "ENDMARK",
+ raw_chunk_plan=plan,
+ )
+ assert_true(report.ok())
+ assert_equal(report.status, 200)
+ assert_equal(report.body_bytes, 600)
+ assert_equal(report.body_match, "yes")
+ assert_equal(report.declared_bytes, 600)
+ assert_equal(report.length_match, "yes")
+ assert_equal(report.surplus_bytes, 0)
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+def test_provider_raw_head_body_reports_malformed_length_grammar() raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "malformed_length", "POST", "/v1/chat/completions", 200, ""
+ )
+ script.raw_response = (
+ "HTTP/1.1 200 OK\r\ncontent-length: 3junk\r\nconnection:"
+ " close\r\n\r\nabc"
+ )
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ var report = run_bounded_call(
+ "raw_head_body",
+ provider_stub.port,
+ 5000,
+ 5000,
+ guard,
+ BOUNDED_CORRELATION_OB_MALFORMED_HEAD,
+ "/v1/chat/completions",
+ "abc",
+ )
+ assert_true(report.completed)
+ assert_true(report.domain_failure())
+ assert_equal(report.cause, "raw_content_length_malformed")
+ provider_stub.wait()
+ assert_true(provider_stub.ok())
+ guard.assert_clean()
+
+
+# OB03: a synthetic write-error seam is never a real peer close.
+
+
+def _assert_synthetic_errno_rejected(errno: Int, declared: Bool) raises:
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "synthetic_errno", "POST", "/v1/chat/completions", 200, '{"choices":[]}'
+ )
+ script.stall_after_head_ms = 200
+ if declared:
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "body_stall"
+ script.inject_write_errno = errno
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ _ = _raw_send_and_read(provider_stub.port, "/v1/chat/completions")
+ provider_stub.reap()
+ assert_true(not provider_stub.ok())
+ assert_equal(provider_stub.phase(), "peer_close")
+ assert_true(provider_stub.reason().find("unexpected_write_") >= 0)
+ assert_true(provider_stub.reason().find("synthdecl") >= 0)
+ guard.assert_clean()
+
+
+def test_provider_scripted_epipe_with_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EPIPE.value), True)
+
+
+def test_provider_scripted_epipe_without_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EPIPE.value), False)
+
+
+def test_provider_scripted_econnreset_with_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.ECONNRESET.value), True)
+
+
+def test_provider_scripted_econnreset_without_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.ECONNRESET.value), False)
+
+
+def test_provider_scripted_eagain_with_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EAGAIN.value), True)
+
+
+def test_provider_scripted_eagain_without_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EAGAIN.value), False)
+
+
+def test_provider_scripted_ebadf_with_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EBADF.value), True)
+
+
+def test_provider_scripted_ebadf_without_declaration_is_not_peer_close() raises:
+ _assert_synthetic_errno_rejected(Int(ErrNo.EBADF.value), False)
+
+
+def _assert_synthetic_errno_delayed_write_rejected(errno: Int) raises:
+ # OB03: the same provenance rule is exercised on the non-stall delayed_write
+ # write step, not only the body_stall step.
+ var scripts = List[ExchangeScript]()
+ var script = exchange_script(
+ "synthetic_delayed",
+ "POST",
+ "/v1/chat/completions",
+ 200,
+ '{"choices":[]}',
+ )
+ script.expect_peer_close = True
+ script.expected_close_cause = "broken_pipe"
+ script.expected_close_phase = "delayed_write"
+ script.inject_write_errno = errno
+ scripts.append(script^)
+ var guard = CleanupGuard()
+ with spawn_max_local_scripted(0, scripts^, guard) as provider_stub:
+ _ = _raw_send_and_read(provider_stub.port, "/v1/chat/completions")
+ provider_stub.reap()
+ assert_true(not provider_stub.ok())
+ assert_equal(provider_stub.phase(), "peer_close")
+ assert_true(provider_stub.reason().find("unexpected_write_") >= 0)
+ assert_true(provider_stub.reason().find("_delayed_write_") >= 0)
+ assert_true(provider_stub.reason().find("synthdecl") >= 0)
+ guard.assert_clean()
+
+
+def test_provider_scripted_epipe_delayed_write_is_not_peer_close() raises:
+ _assert_synthetic_errno_delayed_write_rejected(Int(ErrNo.EPIPE.value))
+
+
+def test_provider_scripted_econnreset_delayed_write_is_not_peer_close() raises:
+ _assert_synthetic_errno_delayed_write_rejected(Int(ErrNo.ECONNRESET.value))
+
+
+def test_provider_scripted_eagain_delayed_write_is_not_peer_close() raises:
+ _assert_synthetic_errno_delayed_write_rejected(Int(ErrNo.EAGAIN.value))
+
+
+def test_provider_scripted_ebadf_delayed_write_is_not_peer_close() raises:
+ _assert_synthetic_errno_delayed_write_rejected(Int(ErrNo.EBADF.value))
+
+
def main() raises:
TestSuite.discover_tests[__functions_in_module()]().run()