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 changes/391b55492c8a85ae68deb9654729681e.yaml
Original file line number Diff line number Diff line change
@@ -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
...
8 changes: 8 additions & 0 deletions changes/c25ec6ffe22011d527e15cec09f6f6be.yaml
Original file line number Diff line number Diff line change
@@ -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
...
7 changes: 7 additions & 0 deletions changes/da926d2af70e7dfa59a5eddb8bbace6c.yaml
Original file line number Diff line number Diff line change
@@ -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
...
1 change: 1 addition & 0 deletions synapse/axon.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
65 changes: 57 additions & 8 deletions synapse/cortex.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import os
import copy
import http
import regex
import asyncio
import logging
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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):
'''
Expand Down Expand Up @@ -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))
Expand All @@ -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

Expand Down
3 changes: 3 additions & 0 deletions synapse/lib/const.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
2 changes: 1 addition & 1 deletion synapse/lib/stormhttp.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
41 changes: 33 additions & 8 deletions synapse/lib/stormtypes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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()
Comment thread
invisig0th marked this conversation as resolved.
return await self.runt.snap.core.axon.metrics()

async def upload(self, genr):
Expand Down Expand Up @@ -2792,18 +2816,19 @@ 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)

sha256 = await tostr(sha256)
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
Expand Down
11 changes: 7 additions & 4 deletions synapse/tests/test_axon.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading
Loading