Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions .github/workflows/decision-trajectories.yml
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,9 @@ name: Private decision trajectories
on:
push:
branches: [master]
paths: [src/bitworld/native_websocket.nim, tests/support/native_websocket_probe.nim, tests/test_native_websocket.py, src/bitworld/decision_trajectory.nim, src/bitworld/runtime.nim, src/bitworld/native_http.nim, src/bitworld/native_stop.nim, src/bitworld/artifact_runtime.nim, tests/test_artifact_http.py, tests/support/artifact_http_probe.nim, tests/test_native_https.py, tests/support/native_https_probe.nim, tests/test_decision_trajectory.nim, tests/test_native_http.py, tests/support/native_http_probe.nim, .github/workflows/decision-trajectories.yml]
paths: [src/bitworld/native_websocket.nim, tests/support/native_websocket_probe.nim, tests/test_native_websocket.py, src/bitworld/decision_trajectory.nim, src/bitworld/runtime.nim, src/bitworld/native_http.nim, src/bitworld/native_stop.nim, src/bitworld/artifact_runtime.nim, tests/test_artifact_http.py, tests/support/artifact_http_probe.nim, tests/test_native_https.py, tests/support/native_https_probe.nim, tests/test_decision_trajectory.nim, tests/test_native_http.py, tests/support/native_http_probe.nim, tests/support/native_request_control_probe.nim, tests/test_native_request_control.py, .github/workflows/decision-trajectories.yml]
pull_request:
paths: [src/bitworld/native_websocket.nim, tests/support/native_websocket_probe.nim, tests/test_native_websocket.py, src/bitworld/decision_trajectory.nim, src/bitworld/runtime.nim, src/bitworld/native_http.nim, src/bitworld/native_stop.nim, src/bitworld/artifact_runtime.nim, tests/test_artifact_http.py, tests/support/artifact_http_probe.nim, tests/test_native_https.py, tests/support/native_https_probe.nim, tests/test_decision_trajectory.nim, tests/test_native_http.py, tests/support/native_http_probe.nim, .github/workflows/decision-trajectories.yml]
paths: [src/bitworld/native_websocket.nim, tests/support/native_websocket_probe.nim, tests/test_native_websocket.py, src/bitworld/decision_trajectory.nim, src/bitworld/runtime.nim, src/bitworld/native_http.nim, src/bitworld/native_stop.nim, src/bitworld/artifact_runtime.nim, tests/test_artifact_http.py, tests/support/artifact_http_probe.nim, tests/test_native_https.py, tests/support/native_https_probe.nim, tests/test_decision_trajectory.nim, tests/test_native_http.py, tests/support/native_http_probe.nim, tests/support/native_request_control_probe.nim, tests/test_native_request_control.py, .github/workflows/decision-trajectories.yml]

jobs:
test:
Expand Down Expand Up @@ -37,6 +37,10 @@ jobs:
run: |
nim c -d:release --threads:on --mm:orc --path:src -o:/tmp/native-http-probe tests/support/native_http_probe.nim
python3 tests/test_native_http.py /tmp/native-http-probe
- name: Verify request-local cancellation and later phase admission
run: |
nim c -d:release --threads:on --mm:orc --path:src -o:/tmp/native-request-control-probe tests/support/native_request_control_probe.nim
python3 tests/test_native_request_control.py /tmp/native-request-control-probe
- name: Verify bounded private artifact finalization after inference interruption
run: |
nim c -d:release --threads:on --mm:orc --path:src -o:/tmp/artifact-http-probe tests/support/artifact_http_probe.nim
Expand Down
8 changes: 6 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -205,9 +205,13 @@ coworld certify among_them/coworld_manifest.json

## Native Coworld HTTP lifecycle

`bitworld/native_http.performNativePost(url, headers, body, deadline)` owns one
`bitworld/native_http.performNativePost(url, headers, body, deadline, control)` owns one
libcurl handle until cleanup. Pass the same absolute `MonoTime` deadline across
attempts. It returns exact received header/body bytes, observed status, actual
attempts. Each request has an owned `NativeRequestControl`; retain it until every
worker using it joins. `cancelNativeRequest(control)` stops only that request and
returns `nhCanceled`. A later decision uses a fresh control without resetting
the original decision deadline. Global stop remains irreversible.
It returns exact received header/body bytes, observed status, actual
transfer completeness, and a typed completion/deadline/interruption/failure kind.
A complete transfer arriving after the deadline remains ineligible for selection.

Expand Down
42 changes: 32 additions & 10 deletions src/bitworld/native_http.nim
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
## No provider parsing or game acceptance occurs here. All received bytes survive
## timeout/interruption; the handle is cleaned before the result can be sealed.

import std/[monotimes, options, os, posix, times]
import std/[atomics, monotimes, options, os, posix, times]
import libcurl except Option
import webby/httpheaders
import native_stop
Expand All @@ -15,7 +15,9 @@ type
RequestPurpose = enum
rpInference, rpArtifact
NativeHttpKind* = enum
nhComplete, nhDeadline, nhInterrupted, nhTransportFailure
nhComplete, nhDeadline, nhInterrupted, nhCanceled, nhTransportFailure
NativeRequestControl* = object
canceled: Atomic[bool]
NativeHttpResponse* = object
kind*: NativeHttpKind
httpStatus*: Option[int]
Expand All @@ -27,6 +29,7 @@ type
Transfer = object
deadline: MonoTime
purpose: RequestPurpose
control: ptr NativeRequestControl
headerBytes, bodyBytes: string

# The pinned Nim binding omits these existing libcurl options.
Expand All @@ -36,6 +39,13 @@ const
OptProtocols = cast[libcurl.Option](181)
OptXferInfoFunction = cast[libcurl.Option](20219)

proc cancelNativeRequest*(control: var NativeRequestControl) {.inline, gcsafe, raises: [].} =
## Irreversible for this request only; the owner retains it until its worker joins.
control.canceled.store(true, moRelaxed)

proc nativeRequestCanceled*(control: var NativeRequestControl): bool {.inline, gcsafe, raises: [].} =
control.canceled.load(moRelaxed)

proc requireCurl(code: Code) =
if code != E_OK:
raise newException(Defect, $easy_strerror(code))
Expand All @@ -62,15 +72,19 @@ proc receiveBody(buffer: cstring, size, count: int, context: pointer): int {.cde
proc checkTransfer(context: pointer, downloadTotal, downloaded,
uploadTotal, uploaded: int64): cint {.cdecl.} =
let transfer = cast[ptr Transfer](context)
if (transfer.purpose == rpInference and interruptionRequested()) or
if (transfer.purpose == rpInference and
(interruptionRequested() or transfer.control[].nativeRequestCanceled())) or
getMonoTime() >= transfer.deadline: 1 else: 0

proc performOwnedRequest(url: string, httpMethod: ArtifactHttpMethod,
headers: HttpHeaders, body: string, deadline: MonoTime,
purpose: RequestPurpose): NativeHttpResponse =
purpose: RequestPurpose, control: var NativeRequestControl): NativeHttpResponse =
if purpose == rpInference and interruptionRequested():
result.kind = nhInterrupted
return
if purpose == rpInference and control.nativeRequestCanceled():
result.kind = nhCanceled
return
let remaining = (deadline - getMonoTime()).inNanoseconds
if remaining <= 0:
result.kind = nhDeadline
Expand All @@ -79,7 +93,7 @@ proc performOwnedRequest(url: string, httpMethod: ArtifactHttpMethod,
let handle = easy_init()
doAssert handle != nil, "Cannot allocate native HTTP handle"
var headerList: Pslist
var transfer = Transfer(deadline: deadline, purpose: purpose)
var transfer = Transfer(deadline: deadline, purpose: purpose, control: control.addr)
var oldMask, pipeMask, previousPending: Sigset
doAssert sigemptyset(pipeMask) == 0
doAssert sigaddset(pipeMask, SIGPIPE) == 0
Expand Down Expand Up @@ -111,8 +125,11 @@ proc performOwnedRequest(url: string, httpMethod: ArtifactHttpMethod,
requireCurl(handle.easy_setopt(OptXferInfoFunction, checkTransfer))
let started = getMonoTime()
let finalRemaining = (deadline - started).inNanoseconds
if finalRemaining <= 0 or (purpose == rpInference and interruptionRequested()):
result.kind = if purpose == rpInference and interruptionRequested(): nhInterrupted else: nhDeadline
if finalRemaining <= 0 or (purpose == rpInference and
(interruptionRequested() or control.nativeRequestCanceled())):
result.kind = if purpose == rpInference and interruptionRequested(): nhInterrupted
elif purpose == rpInference and control.nativeRequestCanceled(): nhCanceled
else: nhDeadline
return
let milliseconds = clong((finalRemaining + 999_999) div 1_000_000)
requireCurl(handle.easy_setopt(OptTimeoutMs, milliseconds))
Expand All @@ -124,6 +141,7 @@ proc performOwnedRequest(url: string, httpMethod: ArtifactHttpMethod,
requireCurl(handle.easy_getinfo(INFO_RESPONSE_CODE, status.addr))
if status != 0: result.httpStatus = some(int(status))
if purpose == rpInference and interruptionRequested(): result.kind = nhInterrupted
elif purpose == rpInference and control.nativeRequestCanceled(): result.kind = nhCanceled
elif code == E_OPERATION_TIMEOUTED or getMonoTime() >= deadline: result.kind = nhDeadline
elif code == E_OK: result.kind = nhComplete
else: result.kind = nhTransportFailure
Expand All @@ -139,17 +157,21 @@ proc performOwnedRequest(url: string, httpMethod: ArtifactHttpMethod,
doAssert sigwait(pipeMask, received) == 0
var discardedMask: Sigset
doAssert pthread_sigmask(SIG_SETMASK, oldMask, discardedMask) == 0
if purpose == rpInference and interruptionRequested(): result.kind = nhInterrupted
elif purpose == rpInference and control.nativeRequestCanceled(): result.kind = nhCanceled
elif getMonoTime() >= deadline: result.kind = nhDeadline
result.headerBytes = move transfer.headerBytes
result.bodyBytes = move transfer.bodyBytes
result.responseReaderJoined = some(true)

proc performNativePost*(url: string, headers: HttpHeaders, body: string,
deadline: MonoTime): NativeHttpResponse =
deadline: MonoTime, control: var NativeRequestControl): NativeHttpResponse =
## A caller shares one deadline across retries. Never reset it per attempt.
performOwnedRequest(url, ahPost, headers, body, deadline, rpInference)
performOwnedRequest(url, ahPost, headers, body, deadline, rpInference, control)

proc performArtifactUpload*(url: string, httpMethod: ArtifactHttpMethod,
headers: HttpHeaders, body: string, cleanupDeadline: MonoTime): NativeHttpResponse =
## Checkpoint finalization has its own finite cleanup lifetime after inference stops.
## No caller can disable interruption in the inference API.
performOwnedRequest(url, httpMethod, headers, body, cleanupDeadline, rpArtifact)
var control: NativeRequestControl
performOwnedRequest(url, httpMethod, headers, body, cleanupDeadline, rpArtifact, control)
3 changes: 2 additions & 1 deletion tests/support/artifact_http_probe.nim
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import std/[base64, json, monotimes, options, os, strutils, times]
import bitworld/[artifact_runtime, decision_trajectory, native_http, native_stop, runtime]

var control: NativeRequestControl
let args = commandLineParams()
installNativeStopHandlers()
let deadline = getMonoTime() + initDuration(milliseconds = args[2].parseBiggestInt())
if args[1] == "stop":
let inference = performNativePost(args[0] & "/inference", @[], "started", deadline)
let inference = performNativePost(args[0] & "/inference", @[], "started", deadline, control)
doAssert inference.kind == nhInterrupted
echo $(%*{"inference": $inference.kind})
flushFile(stdout)
Expand Down
7 changes: 4 additions & 3 deletions tests/support/native_http_probe.nim
Original file line number Diff line number Diff line change
@@ -1,27 +1,28 @@
import std/[base64, json, monotimes, options, os, strutils, times]
import bitworld/[native_http, native_stop]

var control: NativeRequestControl
let args = commandLineParams()
installNativeStopHandlers()
let deadline = getMonoTime() + initDuration(milliseconds = args[1].parseBiggestInt())
let response = performNativePost(args[0], @[("content-type", "application/json")],
"{\"fixture\":true}", deadline)
"{\"fixture\":true}", deadline, control)
doAssert response.responseReaderJoined == some(true)
echo $(%*{"kind": $response.kind, "status": (if response.httpStatus.isSome: %response.httpStatus.get() else: newJNull()),
"headers_b64": encode(response.headerBytes), "body_b64": encode(response.bodyBytes),
"complete": response.transferComplete,
"latency_ms": (if response.latencyMs.isSome: %response.latencyMs.get() else: newJNull())})
if args.len == 3:
let repeated = performNativePost(args[0], @[("content-type", "application/json")],
"{\"fixture\":true}", deadline)
"{\"fixture\":true}", deadline, control)
doAssert repeated.kind == response.kind
doAssert repeated.latencyMs.isNone and repeated.httpStatus.isNone
doAssert repeated.responseReaderJoined.isNone
doAssert repeated.bodyBytes.len == 0 and repeated.headerBytes.len == 0
requestNativeStop()
doAssert interruptionRequested()
let stopped = performNativePost(args[0], @[("content-type", "application/json")],
"{\"fixture\":true}", getMonoTime() + initDuration(seconds = 5))
"{\"fixture\":true}", getMonoTime() + initDuration(seconds = 5), control)
doAssert stopped.kind == nhInterrupted
doAssert stopped.latencyMs.isNone and stopped.httpStatus.isNone
doAssert stopped.bodyBytes.len == 0 and stopped.headerBytes.len == 0
Expand Down
3 changes: 2 additions & 1 deletion tests/support/native_https_probe.nim
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
import std/[json, monotimes, options, os, times]
import bitworld/native_http

var control: NativeRequestControl
let args = commandLineParams()
let deadline = getMonoTime() + initDuration(seconds = 3)
let response = if args[1] == "inference":
performNativePost(args[0], @[], "private fixture", deadline)
performNativePost(args[0], @[], "private fixture", deadline, control)
else:
performArtifactUpload(args[0], ahPut, @[], "private fixture", deadline)
echo $(%*{"kind": $response.kind, "complete": response.transferComplete,
Expand Down
58 changes: 58 additions & 0 deletions tests/support/native_request_control_probe.nim
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
## Parent admits cancellation only after the real HTTP fixture receives a request.
import std/[base64, json, monotimes, options, os, times]
import bitworld/[native_http, native_stop]

type OwnedRequest = object
control: NativeRequestControl
deadline: MonoTime
url: cstring
responseBytes: pointer
responseLen: int

proc request(job: ptr OwnedRequest) {.thread.} =
let response = performNativePost($job.url, @[], "first", job.deadline, job.control)
let resultBytes = $(%*{"kind": $response.kind,
"status": response.httpStatus, "headers_b64": encode(response.headerBytes),
"body_b64": encode(response.bodyBytes), "complete": response.transferComplete,
"reader_joined": response.responseReaderJoined})
job.responseLen = resultBytes.len
job.responseBytes = allocShared(resultBytes.len)
doAssert job.responseBytes != nil
copyMem(job.responseBytes, resultBytes[0].unsafeAddr, resultBytes.len)

let url = commandLineParams()[0]
installNativeStopHandlers()
let deadline = getMonoTime() + initDuration(seconds = 5)
var job = OwnedRequest(deadline: deadline, url: url.cstring)
var worker: Thread[ptr OwnedRequest]
var joined = false
createThread(worker, request, job.addr)
try:
doAssert stdin.readLine() == "cancel"
cancelNativeRequest(job.control)
joinThread(worker)
joined = true
var bytes = newString(job.responseLen)
copyMem(bytes[0].addr, job.responseBytes, bytes.len)
let canceled = parseJson(bytes)
doAssert canceled["kind"].getStr() == "nhCanceled"
doAssert canceled["reader_joined"].getBool()
doAssert not interruptionRequested()
let repeated = performNativePost(url, @[], "not admitted", deadline, job.control)
doAssert repeated.kind == nhCanceled and repeated.httpStatus.isNone
doAssert repeated.responseReaderJoined.isNone
var nextControl: NativeRequestControl
let next = performNativePost(url, @[], "second", deadline, nextControl)
doAssert next.kind == nhComplete and next.responseReaderJoined == some(true)
doAssert next.bodyBytes == "\x00\xffok"
requestNativeStop()
var stoppedControl: NativeRequestControl
let stopped = performNativePost(url, @[], "not admitted", deadline, stoppedControl)
doAssert stopped.kind == nhInterrupted and stopped.httpStatus.isNone
echo $(%*{"canceled": canceled, "later_call_complete": true,
"original_deadline_retained": true, "global_stop_blocks_fresh_control": true})
finally:
if not joined:
cancelNativeRequest(job.control)
joinThread(worker)
if job.responseBytes != nil: deallocShared(job.responseBytes)
64 changes: 64 additions & 0 deletions tests/test_native_request_control.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
"""Real request admission, partial invalid bytes, request-local cancel, and later calls."""
import base64
import http.server
import json
import subprocess
import sys
import threading
import time
from pathlib import Path

probe = Path(sys.argv[1]).resolve()
for trial in range(8):
entered = threading.Event()
release = threading.Event()
requests = []

class Fixture(http.server.BaseHTTPRequestHandler):
def do_POST(self):
body = self.rfile.read(int(self.headers["Content-Length"]))
requests.append(body)
assert body in {b"first", b"second"}
self.send_response(200)
self.send_header("X-Fixture", "first")
self.send_header("X-Fixture", "second")
self.send_header("Content-Length", "4")
self.end_headers()
self.wfile.write(b"\xe2\x82" if body == b"first" else b"\x00\xffok")
self.wfile.flush()
if body == b"first":
entered.set()
release.wait(5)

def log_message(self, *_args):
pass

with http.server.ThreadingHTTPServer(("127.0.0.1", 0), Fixture) as server:
owner = threading.Thread(target=server.serve_forever)
owner.start()
process = subprocess.Popen([str(probe), f"http://127.0.0.1:{server.server_port}/v1/messages"],
stdin=subprocess.PIPE, stdout=subprocess.PIPE,
stderr=subprocess.PIPE, text=True)
try:
assert entered.wait(3), "the real owned request never entered"
time.sleep(0.2)
started = time.monotonic()
stdout, stderr = process.communicate("cancel\n", timeout=3)
assert process.returncode == 0, stderr
result = json.loads(stdout)
assert result["canceled"]["status"] == 200
assert not result["canceled"]["complete"]
assert base64.b64decode(result["canceled"]["body_b64"], validate=True) == b"\xe2\x82"
assert b"X-Fixture: first\r\nX-Fixture: second\r\n" in base64.b64decode(result["canceled"]["headers_b64"], validate=True)
assert result["later_call_complete"] and result["original_deadline_retained"]
assert result["global_stop_blocks_fresh_control"]
assert requests == [b"first", b"second"]
print(trial, "phase cancel joined; later call completes; global stop seals", time.monotonic() - started, flush=True)
finally:
if process.poll() is None:
process.terminate()
process.wait(timeout=3)
release.set()
server.shutdown()
owner.join(timeout=3)
assert not owner.is_alive()
Loading