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
7 changes: 7 additions & 0 deletions docs/website/src/impl-spec/appendix/manager-socket.rst
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,13 @@ and process spawn continue asynchronously and report through ``event``
notifications. The requesting connection is subscribed to the run's events
automatically.

Only decoding is synchronous, because the request has to be read before an id
can be attached to it: calldata that fails to decode, or a payload shape that
names no ``run``, is answered with ``error``. Everything decided from the
decoded request is asynchronous and arrives as a terminal event. A non-empty
``host_hello_data[1]``, which is manager-owned, and a request that needs
modules while they are stopped, both report ``failed_to_start`` this way.

``attach``
----------

Expand Down
34 changes: 6 additions & 28 deletions implementation/src/manager/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -491,35 +491,13 @@ where
return Ok(());
}
};
if req
.host_hello_data
.get(1)
.is_some_and(|data| !data.is_empty())
{
self.writer
.error(
request_id,
Errors::MalformedFrame,
"host_hello_data for manager-owned host index 1 is not allowed",
)
.await?;
return Ok(());
}

// `host_hello_data[1]` and a modules-stopped request are rejected by
// supervision (`run_genvm_process`), not here: the protocol returns the
// genvm_id immediately and reports those failures as a terminal event.
// Only the module locks are taken on this path, to pin the handlers
// before the run spawns.
let modules_lock = if req.needs_modules() {
match super::modules::Ctx::get_module_locks(self.ctx.gep(|x| &x.mod_ctx)).await {
Some(lock) => Some(lock),
None => {
self.writer
.error(
request_id,
Errors::Internal,
"modules are required but not running",
)
.await?;
return Ok(());
}
}
super::modules::Ctx::get_module_locks(self.ctx.gep(|x| &x.mod_ctx)).await
} else {
None
};
Expand Down
47 changes: 47 additions & 0 deletions tests/system/manager-socket/test.py
Original file line number Diff line number Diff line change
Expand Up @@ -402,6 +402,46 @@ async def _oversized_closes_connection(self):
aiohttp.WSCloseCode.MESSAGE_TOO_BIG,
), close_code

async def _run_returns_id_before_semantic_rejection(self):
"""
Both semantic rejections the socket used to answer synchronously must now
come back as a terminal event, after the genvm_id has been returned.
"""
async with await self._manager('async-validation') as manager:
host_path = manager.work_dir / 'unused-host.sock'

# `host_hello_data[1]` is manager-owned, so a non-empty slot is a
# semantic rejection rather than a malformed frame.
async with ManagerWsClient(manager.uri) as client:
await _read_hello(client)
req = _run_request(self.case.shared.root_dir, host_path)
req['run']['host_hello_data'] = [b'', b'injected']
await client.send(Methods.RUN, 10, req)
method, request_id, payload = await client.read_frame()
assert (method, request_id) == (Methods.RUN, 10)
genvm_id = payload['genvm_id']
variant, event = await _wait_terminal(client, timeout=5)
assert variant == 'failed_to_start', variant
assert event['genvm_id'] == genvm_id
assert 'host index 1' in event['error'], event

# A request needing modules with the modules stopped is the other
# class: it used to answer with an error and no genvm_id.
async with ManagerWsClient(manager.uri) as client:
await _read_hello(client)
req = _run_request(self.case.shared.root_dir, host_path)
req['run']['no_modules'] = False
req['run']['is_sync'] = False
req['run']['permissions'] = 'n'
await client.send(Methods.RUN, 20, req)
method, request_id, payload = await client.read_frame()
assert (method, request_id) == (Methods.RUN, 20)
genvm_id = payload['genvm_id']
variant, event = await _wait_terminal(client, timeout=5)
assert variant == 'failed_to_start', variant
assert event['genvm_id'] == genvm_id
assert 'modules are required' in event['error'], event

async def _startup_failure_events_and_permits(self):
async with await self._manager('startup-failures') as manager:
async with ManagerWsClient(manager.uri) as client:
Expand Down Expand Up @@ -958,6 +998,13 @@ def service(
None,
genvm.ManagerService,
),
(
'async-validation',
'_run_returns_id_before_semantic_rejection',
{},
None,
genvm.ManagerService,
),
('happy-path', '_happy_path_artifact_ack_attach', {}, None, genvm.ManagerService),
(
'reconnect-mid-run',
Expand Down
Loading