diff --git a/changes/391b55492c8a85ae68deb9654729681e.yaml b/changes/391b55492c8a85ae68deb9654729681e.yaml new file mode 100644 index 0000000000..15c36a1444 --- /dev/null +++ b/changes/391b55492c8a85ae68deb9654729681e.yaml @@ -0,0 +1,7 @@ +--- +desc: Fixed the Axon HTTP upload API starting an upload in the Axon for a request + which was denied. +desc:literal: false +prs: [] +type: bug +... diff --git a/changes/c25ec6ffe22011d527e15cec09f6f6be.yaml b/changes/c25ec6ffe22011d527e15cec09f6f6be.yaml new file mode 100644 index 0000000000..073861ec54 --- /dev/null +++ b/changes/c25ec6ffe22011d527e15cec09f6f6be.yaml @@ -0,0 +1,8 @@ +--- +desc: Fixed a ``KeyError`` in Storm APIs which interacted with the Axon. Cortex Storm, + Telepath, and HTTP APIs which use the Axon now wait up to 300 seconds for the Axon to be + ready before raising a ``TimeOut`` error. +desc:literal: false +prs: [] +type: bug +... diff --git a/changes/da926d2af70e7dfa59a5eddb8bbace6c.yaml b/changes/da926d2af70e7dfa59a5eddb8bbace6c.yaml new file mode 100644 index 0000000000..d16fa4f7ba --- /dev/null +++ b/changes/da926d2af70e7dfa59a5eddb8bbace6c.yaml @@ -0,0 +1,7 @@ +--- +desc: Added an ``axon:ready`` key to the ``cell`` section of the Cortex ``getCellInfo()`` + API, which reports whether the Axon is ready. +desc:literal: false +prs: [] +type: feat +... diff --git a/synapse/axon.py b/synapse/axon.py index fc6495057b..f7c3f54530 100644 --- a/synapse/axon.py +++ b/synapse/axon.py @@ -48,6 +48,7 @@ async def prepare(self): if not await self.allowed(('axon', 'upload')): await self.finish() + return # max_body_size defaults to 100MB and requires a value self.request.connection.set_max_body_size(MAX_HTTP_UPLOAD_SIZE) diff --git a/synapse/cortex.py b/synapse/cortex.py index 7bc9185842..059d76be79 100644 --- a/synapse/cortex.py +++ b/synapse/cortex.py @@ -1,5 +1,6 @@ import os import copy +import http import regex import asyncio import logging @@ -175,9 +176,18 @@ async def wrap_liftgenr(iden, genr): class CortexAxonMixin: - async def prepare(self): - await self.cell.axready.wait() - await s_coro.ornot(super().prepare) + async def allowed(self, perm, default=False, gateiden=None): + # wait for the Axon only once the user has the permission + if not await super().allowed(perm, default=default, gateiden=gateiden): + return False + + try: + await self.cell.getAxon() + except s_exc.TimeOut as e: + self.sendRestExc(e, status_code=http.HTTPStatus.SERVICE_UNAVAILABLE) + return False + + return True def getAxon(self): return self.cell.axon @@ -760,13 +770,13 @@ async def iterTagPropRows(self, layriden, tag, prop, form=None, stortype=None, s async def getAxonUpload(self): self.user.confirm(('axon', 'upload')) - await self.cell.axready.wait() + await self.cell.getAxon() upload = await self.cell.axon.upload() return await s_axon.UpLoadProxy.anit(self.link, upload) async def getAxonBytes(self, sha256): self.user.confirm(('axon', 'get')) - await self.cell.axready.wait() + await self.cell.getAxon() async for byts in self.cell.axon.get(s_common.uhex(sha256)): yield byts @@ -6168,9 +6178,45 @@ def getStormLib(self, path): def getStormCmds(self): return list(self.stormcmds.items()) - async def getAxon(self): - await self.axready.wait() - return self.axon.iden + async def getAxon(self, timeout=s_const.AXON_READY_TIMEOUT): + ''' + Wait for the Axon to be ready and return its iden. + + Args: + timeout (int): The maximum number of seconds to wait, or None to wait indefinitely. + + Returns: + str: The iden of the Axon, or None if the Axon did not report one. + + Raises: + s_exc.TimeOut: If the Axon is not ready within the timeout. + ''' + if not self.axready.is_set(): + try: + await s_common.wait_for(self.axready.wait(), timeout) + except asyncio.TimeoutError: + mesg = f'Timed out waiting {timeout} seconds for the Axon to be ready.' + raise s_exc.TimeOut(mesg=mesg, timeout=timeout) from None + + if (cellinfo := self.axoninfo.get('cell')) is None: + return None + + return cellinfo.get('iden') + + async def getCellInfo(self): + ''' + Return metadata specific for the Cortex. + + Notes: + In addition to the base Cell information, the ``cell`` section + includes ``axon:ready``, which is True when the Axon is ready. + + Returns: + Dict: A Dictionary of metadata. + ''' + info = await super().getCellInfo() + info['cell']['axon:ready'] = self.axready.is_set() + return info def setFeedFunc(self, name, func): ''' @@ -6493,6 +6539,7 @@ async def exportStorm(self, text, opts=None): yield pode async def exportStormToAxon(self, text, opts=None): + await self.getAxon() async with await self.axon.upload() as fd: async for pode in self.exportStorm(text, opts=opts): await fd.write(s_msgpack.en(pode)) @@ -6513,6 +6560,8 @@ async def feedFromAxon(self, sha256, opts=None): # ensure that the user can make all node edits in the layer user.confirm(('node',), gateiden=view.layers[0].iden) + await self.getAxon() + q = s_queue.Queue(maxsize=10000) feedexc = None diff --git a/synapse/lib/const.py b/synapse/lib/const.py index f049ccf985..71414f41c3 100644 --- a/synapse/lib/const.py +++ b/synapse/lib/const.py @@ -52,5 +52,8 @@ MAX_LINE_SIZE = kibibyte * 64 MAX_FIELD_SIZE = kibibyte * 64 +# Axon constants +AXON_READY_TIMEOUT = 300 # seconds + # Socket constants UNIX_SOCKET_PATH_MAX = 103 diff --git a/synapse/lib/stormhttp.py b/synapse/lib/stormhttp.py index 40610f6240..ac7a31ebc0 100644 --- a/synapse/lib/stormhttp.py +++ b/synapse/lib/stormhttp.py @@ -425,7 +425,7 @@ async def _httpRequest(self, meth, url, headers=None, json=None, body=None, kwargs['proxy'] = proxy if ssl_opts is not None: - axonvers = self.runt.snap.core.axoninfo['synapse']['version'] + axonvers = await s_stormtypes.getAxonSynapseVersion(self.runt) mesg = f'The ssl_opts argument requires an Axon Synapse version {s_stormtypes.AXON_MINVERS_SSLOPTS}, ' \ f'but the Axon is running {axonvers}' s_version.reqVersion(axonvers, s_stormtypes.AXON_MINVERS_SSLOPTS, mesg=mesg) diff --git a/synapse/lib/stormtypes.py b/synapse/lib/stormtypes.py index a5f1751d8a..9e22a05a82 100644 --- a/synapse/lib/stormtypes.py +++ b/synapse/lib/stormtypes.py @@ -108,6 +108,28 @@ async def resolveCoreProxyUrl(valu): case _: raise s_exc.BadArg(mesg='HTTP proxy argument must be a string or bool.') +async def getAxonSynapseVersion(runt): + ''' + Get the Synapse version of the Cortex's Axon. + + Args: + runt (Runtime): The Storm runtime. + + Returns: + tuple: The Synapse version of the Axon. + + Raises: + s_exc.FeatureNotSupported: If the Axon version is unavailable. + s_exc.TimeOut: If the Axon is not ready within the default Axon timeout. + ''' + await runt.snap.core.getAxon() + + if (syninfo := runt.snap.core.axoninfo.get('synapse')) is None or (axonvers := syninfo.get('version')) is None: + mesg = 'Unable to determine the Synapse version of the Axon.' + raise s_exc.FeatureNotSupported(mesg=mesg) + + return axonvers + async def resolveAxonProxyArg(valu): ''' Resolve a proxy value to the kwarg to set for an Axon HTTP call. @@ -120,7 +142,7 @@ async def resolveAxonProxyArg(valu): ''' runt = s_scope.get('runt') - axonvers = runt.snap.core.axoninfo['synapse']['version'] + axonvers = await getAxonSynapseVersion(runt) if axonvers < AXON_MINVERS_PROXY: await runt.snap.warnonce(f'Axon version does not support proxy argument: {axonvers} < {AXON_MINVERS_PROXY}') return False, None @@ -2556,7 +2578,7 @@ async def wget(self, url, headers=None, params=None, method='GET', json=None, bo kwargs['proxy'] = proxy if ssl_opts is not None: - axonvers = self.runt.snap.core.axoninfo['synapse']['version'] + axonvers = await getAxonSynapseVersion(self.runt) mesg = f'The ssl_opts argument requires an Axon Synapse version {AXON_MINVERS_SSLOPTS}, ' \ f'but the Axon is running {axonvers}' s_version.reqVersion(axonvers, AXON_MINVERS_SSLOPTS, mesg=mesg) @@ -2597,7 +2619,7 @@ async def wput(self, sha256, url, headers=None, params=None, method='PUT', kwargs['proxy'] = proxy if ssl_opts is not None: - axonvers = self.runt.snap.core.axoninfo['synapse']['version'] + axonvers = await getAxonSynapseVersion(self.runt) mesg = f'The ssl_opts argument requires an Axon Synapse version {AXON_MINVERS_SSLOPTS}, ' \ f'but the Axon is running {axonvers}' s_version.reqVersion(axonvers, AXON_MINVERS_SSLOPTS, mesg=mesg) @@ -2706,6 +2728,8 @@ async def csvrows(self, sha256, dialect='excel', errors='ignore', **fmtparams): async def metrics(self): if not self.runt.allowed(('axon', 'has')): self.runt.confirm(('storm', 'lib', 'axon', 'has')) + + await self.runt.snap.core.getAxon() return await self.runt.snap.core.axon.metrics() async def upload(self, genr): @@ -2792,7 +2816,12 @@ async def unpack(self, sha256, fmt, offs=0): ''' Unpack bytes from a file in the Axon using struct. ''' - if self.runt.snap.core.axoninfo.get('features', {}).get('unpack', 0) < 1: + if not self.runt.allowed(('axon', 'get')): + self.runt.confirm(('storm', 'lib', 'axon', 'get')) + + await self.runt.snap.core.getAxon() + + if (features := self.runt.snap.core.axoninfo.get('features')) is None or features.get('unpack', 0) < 1: mesg = 'The connected Axon does not support the the unpack API. Please update your Axon.' raise s_exc.FeatureNotSupported(mesg=mesg) @@ -2800,10 +2829,6 @@ async def unpack(self, sha256, fmt, offs=0): fmt = await tostr(fmt) offs = await toint(offs) - if not self.runt.allowed(('axon', 'get')): - self.runt.confirm(('storm', 'lib', 'axon', 'get')) - - await self.runt.snap.core.getAxon() return await self.runt.snap.core.axon.unpack(s_common.uhex(sha256), fmt, offs) @registry.registerLib diff --git a/synapse/tests/test_axon.py b/synapse/tests/test_axon.py index e5737d8e90..0906089edd 100644 --- a/synapse/tests/test_axon.py +++ b/synapse/tests/test_axon.py @@ -612,10 +612,13 @@ async def runAxonTestHttp(self, axon, realaxon=None): item = await resp.json() self.eq('err', item.get('status')) - async with sess.post(url_ul, data=abuf) as resp: - self.eq(resp.status, http.HTTPStatus.FORBIDDEN) - item = await resp.json() - self.eq('err', item.get('status')) + # a denied upload does not start an upload in the Axon + with mock.patch.object(realaxon, 'upload', wraps=realaxon.upload) as upload: + async with sess.post(url_ul, data=abuf) as resp: + self.eq(resp.status, http.HTTPStatus.FORBIDDEN) + item = await resp.json() + self.eq('err', item.get('status')) + upload.assert_not_called() # Stream file byts = io.BytesIO(bbuf) diff --git a/synapse/tests/test_cortex.py b/synapse/tests/test_cortex.py index 74083c7f03..56e28ad6f9 100644 --- a/synapse/tests/test_cortex.py +++ b/synapse/tests/test_cortex.py @@ -4,7 +4,9 @@ import time import asyncio import hashlib +import inspect import logging +import functools import regex @@ -20,6 +22,7 @@ import synapse.lib.coro as s_coro import synapse.lib.node as s_node import synapse.lib.time as s_time +import synapse.lib.const as s_const import synapse.lib.layer as s_layer import synapse.lib.storm as s_storm import synapse.lib.output as s_output @@ -6863,6 +6866,11 @@ async def test_cortex_axon(self): self.eq(size, 8) self.eq(s_common.ehex(sha2), '2413fb3709b05939f04cf2e92f7d0897fc2596f9ad0b8a9ea855c7bfebaae892') self.true(core.nexsroot is core.axon.nexsroot) + self.eq(await core.getAxon(), core.axon.iden) + + info = await core.getCellInfo() + self.true(info['cell']['axon:ready']) + self.notin('axon:version', info['cell']) self.true(core.axon.isfini) self.false(core.axready.is_set()) @@ -6900,6 +6908,134 @@ async def test_cortex_axon(self): self.eq(await axon.metrics(), await core.axon.metrics()) + async def test_cortex_axon_ready_timeout(self): + + with self.getTestDir() as dirn: + + async with self.getTestAxon(dirn=dirn) as axon: + aurl = axon.getLocalUrl() + + async with self.getTestCore(conf={'axon': aurl}) as core: + + self.false(core.axready.is_set()) + self.eq(core.axoninfo, {}) + + info = await core.getCellInfo() + self.false(info['cell']['axon:ready']) + + sha256 = s_common.ehex(hashlib.sha256(b'vertex').digest()) + opts = {'vars': {'sha256': sha256, 'url': 'http://127.0.0.1:1/'}} + queries = ( + 'return($lib.axon.wget($url))', + 'return($lib.axon.wput($sha256, $url))', + 'return($lib.axon.unpack($sha256, fmt=">Q"))', + 'yield $lib.axon.urlfile($url)', + ''' + $fields = ([{"name": "file", "sha256": $sha256}]) + return($lib.inet.http.post($url, fields=$fields)) + ''', + 'for $line in $lib.axon.readlines($sha256) {}', + 'for $item in $lib.axon.jsonlines($sha256) {}', + 'for $row in $lib.axon.csvrows($sha256) {}', + 'for $item in $lib.axon.list() {}', + 'return($lib.axon.dels(($sha256,)))', + 'return($lib.axon.del($sha256))', + 'return($lib.axon.upload(([])))', + 'return($lib.axon.has($sha256))', + 'return($lib.axon.size($sha256))', + 'return($lib.axon.put($buf))', + 'return($lib.axon.hashset($sha256))', + 'return($lib.axon.read($sha256))', + 'return($lib.axon.metrics())', + 'return($lib.feed.fromAxon($sha256))', + 'return($lib.bytes.put($buf))', + 'return($lib.bytes.has($sha256))', + 'return($lib.bytes.size($sha256))', + 'return($lib.bytes.hashset($sha256))', + 'return($lib.bytes.upload(([])))', + ) + opts['vars']['buf'] = b'vertex' + + # $lib.bytes permissions are allowed by default, so deny them explicitly + visi = await core.auth.addUser('visi') + await visi.addRule((False, ('axon',))) + visiopts = {'user': visi.iden, 'vars': opts['vars']} + + timeout = inspect.signature(core.getAxon).parameters['timeout'].default + self.eq(timeout, s_const.AXON_READY_TIMEOUT) + + with self.raises(s_exc.TimeOut) as cm: + await core.getAxon(timeout=0.1) + self.eq(cm.exception.get('mesg'), 'Timed out waiting 0.1 seconds for the Axon to be ready.') + self.eq(cm.exception.get('timeout'), 0.1) + + with patch.object(core, 'getAxon', functools.partial(core.getAxon, timeout=0.1)): + + async with core.getLocalProxy() as proxy: + + with self.raises(s_exc.TimeOut): + await proxy.getAxonUpload() + + with self.raises(s_exc.TimeOut): + async for byts in proxy.getAxonBytes(sha256): + pass + + with self.raises(s_exc.TimeOut): + await proxy.feedFromAxon(sha256) + + for query in queries: + with self.raises(s_exc.TimeOut): + await core.callStorm(query, opts=opts) + + # permission checks run before waiting on the Axon + for query in queries: + with self.raises(s_exc.AuthDeny): + await core.callStorm(query, opts=visiopts) + + # an export has no permission of its own to check first + with self.raises(s_exc.TimeOut): + await core.callStorm('return($lib.export.toaxon("inet:fqdn"))') + + async with self.getTestAxon(dirn=dirn) as axon: + + self.true(await s_coro.event_wait(core.axready, timeout=10)) + self.nn(core.axoninfo['synapse']['version']) + + info = await core.callStorm('return($lib.cell.getCellInfo())') + self.true(info['cell']['axon:ready']) + + await core.axon.put(b'vertex') + + self.eq(await core.getAxon(), axon.iden) + self.eq(await core.getAxon(timeout=None), axon.iden) + + originfo = core.axoninfo + try: + for axoninfo in ({}, {'cell': {}}): + core.axoninfo = axoninfo + self.none(await core.getAxon()) + + finally: + core.axoninfo = originfo + + resp = await core.callStorm(queries[0], opts=opts) + self.false(resp.get('ok')) + + resp = await core.callStorm(queries[1], opts=opts) + self.false(resp.get('ok')) + + resp = await core.callStorm(queries[4], opts=opts) + self.eq(resp.get('code'), -1) + + # the Axon is reported as not ready after it disconnects + for _ in range(20): + if not core.axready.is_set(): + break + await asyncio.sleep(0.1) + + info = await core.getCellInfo() + self.false(info['cell']['axon:ready']) + async def test_cortex_delLayerView(self): with self.getTestDir() as dirn: diff --git a/synapse/tests/test_lib_httpapi.py b/synapse/tests/test_lib_httpapi.py index 119b50c282..ead0ab7d6e 100644 --- a/synapse/tests/test_lib_httpapi.py +++ b/synapse/tests/test_lib_httpapi.py @@ -1,5 +1,8 @@ import ssl import http +import functools + +from unittest import mock import aiohttp import aiohttp.client_exceptions as a_exc @@ -2193,15 +2196,55 @@ async def test_core_remote_axon_http(self): host, port = await core.addHttpsPort(0, host='127.0.0.1') + root = await core.auth.getUserByName('root') + await root.setPasswd('root') + + newb = await core.auth.addUser('newb') + await newb.setPasswd('secret') + + await axon.fini() + + sha256 = s_common.ehex(s_t_axon.asdfhash) + hasurl = f'https://localhost:{port}/api/v1/axon/files/has/sha256/{sha256}' + puturl = f'https://localhost:{port}/api/v1/axon/files/put' + + # auth and permission checks run before waiting on the Axon async with self.getHttpSess() as sess: - await axon.fini() + async with sess.get(hasurl, timeout=timeout) as resp: + self.eq(resp.status, http.HTTPStatus.UNAUTHORIZED) + + async with self.getHttpSess(auth=('newb', 'secret'), port=port) as sess: + + async with sess.get(hasurl, timeout=timeout) as resp: + self.eq(resp.status, http.HTTPStatus.FORBIDDEN) + + async with sess.post(puturl, data=b'asdfasdf', timeout=timeout) as resp: + self.eq(resp.status, http.HTTPStatus.FORBIDDEN) + + async with self.getHttpSess(auth=('root', 'root'), port=port) as sess: with self.raises(TimeoutError): - sha256 = s_common.ehex(s_t_axon.asdfhash) - url = f'https://localhost:{port}/api/v1/axon/files/has/sha256/{sha256}' - async with sess.get(url, timeout=timeout) as resp: + async with sess.get(hasurl, timeout=timeout) as resp: pass + with mock.patch.object(core, 'getAxon', functools.partial(core.getAxon, timeout=0.1)): + + async with sess.get(hasurl) as resp: + self.eq(resp.status, http.HTTPStatus.SERVICE_UNAVAILABLE) + item = await resp.json() + self.eq(item.get('status'), 'err') + self.eq(item.get('code'), 'TimeOut') + self.eq(item.get('mesg'), 'Timed out waiting 0.1 seconds for the Axon to be ready.') + + # the upload handler finishes without starting an upload + with self.getLoggerStream('tornado.application') as stream: + async with sess.post(puturl, data=b'asdfasdf') as resp: + self.eq(resp.status, http.HTTPStatus.SERVICE_UNAVAILABLE) + item = await resp.json() + self.eq(item.get('code'), 'TimeOut') + + self.notin('Uncaught exception', stream.getvalue()) + async def test_http_login_broken(self): async with self.getTestCore() as core: diff --git a/synapse/tests/test_lib_storm.py b/synapse/tests/test_lib_storm.py index 2e4056ce3f..71fd0e50cd 100644 --- a/synapse/tests/test_lib_storm.py +++ b/synapse/tests/test_lib_storm.py @@ -1,6 +1,7 @@ import copy import asyncio import textwrap +import functools import itertools import urllib.parse as u_parse import unittest.mock as mock @@ -6289,7 +6290,23 @@ async def test_lib_storm_delnode(self): with self.raises(s_exc.AuthDeny): await asvisi.callStorm(f'file:bytes={sha256} | delnode --delbytes') - await visi.addRule((True, ('storm', 'lib', 'axon', 'del'))) + core.axready.clear() + try: + with mock.patch.object(core, 'getAxon', functools.partial(core.getAxon, timeout=0.1)): + + # permission checks run before waiting on the Axon + with self.raises(s_exc.AuthDeny): + await asvisi.callStorm(f'file:bytes={sha256} | delnode --delbytes') + + await visi.addRule((True, ('storm', 'lib', 'axon', 'del'))) + + with self.raises(s_exc.TimeOut): + await asvisi.callStorm(f'file:bytes={sha256} | delnode --delbytes') + + self.len(1, await core.nodes(f'file:bytes={sha256}')) + + finally: + core.axready.set() await asvisi.callStorm(f'file:bytes={sha256} | delnode --delbytes') self.len(0, await core.nodes(f'file:bytes={sha256}')) diff --git a/synapse/tests/test_lib_stormhttp.py b/synapse/tests/test_lib_stormhttp.py index 3f8f16deea..80d9d3ed31 100644 --- a/synapse/tests/test_lib_stormhttp.py +++ b/synapse/tests/test_lib_stormhttp.py @@ -662,6 +662,40 @@ async def test_storm_http_post_file(self): self.eq(code, -1) self.eq('ValueError', errname) + async def test_storm_http_post_file_axon_info_missing(self): + + async with self.getTestCore() as core: + + size, sha256 = await core.axon.put(b'vertex') + opts = {'vars': {'sha256': s_common.ehex(sha256), 'url': 'http://127.0.0.1:1/'}} + + queries = ( + ''' + $fields = ([{"name": "file", "sha256": $sha256}]) + return($lib.inet.http.post($url, fields=$fields)) + ''', + ''' + $fields = ([{"name": "file", "sha256": $sha256}]) + return($lib.inet.http.post($url, fields=$fields, ssl_opts=({"verify": (false)}))) + ''', + ) + + originfo = core.axoninfo + try: + for axoninfo in ({}, {'synapse': {}}): + core.axoninfo = axoninfo + for query in queries: + with self.raises(s_exc.FeatureNotSupported) as cm: + await core.callStorm(query, opts=opts) + self.eq(cm.exception.get('mesg'), 'Unable to determine the Synapse version of the Axon.') + + finally: + core.axoninfo = originfo + + for query in queries: + resp = await core.callStorm(query, opts=opts) + self.eq(resp.get('code'), -1) + async def test_storm_http_proxy(self): conf = {'http:proxy': 'socks5://user:pass@127.0.0.1:1'} async with self.getTestCore(conf=conf) as core: diff --git a/synapse/tests/test_lib_stormlib_imap.py b/synapse/tests/test_lib_stormlib_imap.py index 6160a9b01e..66f18b69e4 100644 --- a/synapse/tests/test_lib_stormlib_imap.py +++ b/synapse/tests/test_lib_stormlib_imap.py @@ -3,6 +3,7 @@ import imaplib import logging import textwrap +import functools import contextlib import regex @@ -883,6 +884,32 @@ async def test_storm_imap_fetch(self): self.eq(ret, (True, (rfc822, header, b'(UID 1 RFC822 BODY[HEADER])'))) + async def test_storm_imap_fetch_axon_timeout(self): + + async with self.getTestCoreAndImapPort() as (core, port): + user = 'user00@vertex.link' + opts = {'vars': {'port': port, 'user': user}} + + scmd = ''' + $server = $lib.inet.imap.connect(127.0.0.1, port=$port, ssl=(false)) + $server.login($user, "pass00") + $server.select("INBOX") + yield $server.fetch("1") + ''' + + core.axready.clear() + try: + with mock.patch.object(core, 'getAxon', functools.partial(core.getAxon, timeout=0.1)): + with self.raises(s_exc.TimeOut): + await core.nodes(scmd, opts=opts) + + finally: + core.axready.set() + + nodes = await core.nodes(scmd, opts=opts) + self.len(1, nodes) + self.eq('file:bytes', nodes[0].ndef[0]) + async def test_storm_imap_logout(self): async with self.getTestCoreAndImapPort() as (core, port): diff --git a/synapse/tests/test_lib_stormtypes.py b/synapse/tests/test_lib_stormtypes.py index 1c57b8b04d..31a6aee6a4 100644 --- a/synapse/tests/test_lib_stormtypes.py +++ b/synapse/tests/test_lib_stormtypes.py @@ -8199,6 +8199,38 @@ async def fake(): with self.raises(s_exc.BadState): await core.callStorm(merging) + async def test_storm_lib_axon_info_missing(self): + + async with self.getTestCore() as core: + + size, sha256 = await core.axon.put(b'vertex') + opts = {'vars': {'sha256': s_common.ehex(sha256), 'url': 'http://127.0.0.1:1/'}} + + queries = ( + 'return($lib.axon.wget($url))', + 'return($lib.axon.wget($url, proxy=(false)))', + 'return($lib.axon.wget($url, ssl_opts=({"verify": (false)})))', + 'return($lib.axon.wput($sha256, $url))', + 'return($lib.axon.wput($sha256, $url, ssl_opts=({"verify": (false)})))', + 'yield $lib.axon.urlfile($url)', + ) + + originfo = core.axoninfo + try: + for axoninfo in ({}, {'synapse': {}}): + core.axoninfo = axoninfo + for query in queries: + with self.raises(s_exc.FeatureNotSupported) as cm: + await core.callStorm(query, opts=opts) + self.eq(cm.exception.get('mesg'), 'Unable to determine the Synapse version of the Axon.') + + finally: + core.axoninfo = originfo + + for query in queries[:-1]: + resp = await core.callStorm(query, opts=opts) + self.false(resp.get('ok')) + async def test_storm_lib_axon_read_unpack(self): async with self.getTestCore() as core: @@ -8206,14 +8238,18 @@ async def test_storm_lib_axon_read_unpack(self): visi = await core.auth.addUser('visi') orig_axoninfo = core.axoninfo - core.axoninfo = {'features': {}} data = struct.pack('>Q', 1) size, sha256 = await core.axon.put(data) sha256_s = s_common.ehex(sha256) q = 'return($lib.axon.unpack($sha256, fmt=">Q"))' - await self.asyncraises(s_exc.FeatureNotSupported, - core.callStorm(q, opts={'vars': {'sha256': sha256_s}})) - core.axoninfo = orig_axoninfo + try: + for axoninfo in ({}, {'features': {}}): + core.axoninfo = axoninfo + with self.raises(s_exc.FeatureNotSupported): + await core.callStorm(q, opts={'vars': {'sha256': sha256_s}}) + + finally: + core.axoninfo = orig_axoninfo data = b'vertex.link' size, sha256 = await core.axon.put(data)