82 lines
3.4 KiB
Python
82 lines
3.4 KiB
Python
"""Local ClamAV boundary. Unknown/error/over-limit results NEVER release artwork."""
|
|
import socket
|
|
import struct
|
|
import time
|
|
from fastapi import HTTPException
|
|
from .core.auth import audit
|
|
from .core.db import connect
|
|
from .core.limits import scan_limit_bytes
|
|
|
|
def require_clean(row):
|
|
if not row['complete'] or row['scan_state'] != 'clean':
|
|
raise HTTPException(409, 'Artwork is quarantined until the malware scan succeeds')
|
|
|
|
class ClamAV:
|
|
def command(self, command, timeout=2):
|
|
with socket.create_connection(('scanner',3310),timeout=timeout) as sock:
|
|
sock.settimeout(timeout)
|
|
sock.sendall(b'z'+command+b'\0')
|
|
reply=b''
|
|
while b'\0' not in reply and len(reply)<4096:
|
|
block=sock.recv(4096)
|
|
if not block:break
|
|
reply+=block
|
|
return reply.rstrip(b'\0\n')
|
|
|
|
def ping(self):
|
|
return self.command(b'PING') == b'PONG'
|
|
|
|
def version(self):
|
|
return self.command(b'VERSION').decode('utf-8','replace')
|
|
|
|
def scan(self, stream, size):
|
|
if size > scan_limit_bytes():
|
|
return 'rejected', 'File exceeds the malware scan limit'
|
|
with socket.create_connection(('scanner',3310),timeout=10) as sock:
|
|
sock.settimeout(150)
|
|
sock.sendall(b'zINSTREAM\0')
|
|
sent=0
|
|
for chunk in stream.iter_chunks(chunk_size=65536):
|
|
sent+=len(chunk)
|
|
if sent>size: return 'rejected','Stored file size changed'
|
|
sock.sendall(struct.pack('!I',len(chunk))+chunk)
|
|
if sent!=size:return 'rejected','Stored file size changed'
|
|
sock.sendall(b'\0\0\0\0')
|
|
reply=b''
|
|
while b'\0' not in reply and len(reply)<4096:
|
|
block=sock.recv(4096)
|
|
if not block:break
|
|
reply+=block
|
|
result=reply.rstrip(b'\0\n')
|
|
if result==b'stream: OK':return 'clean',None
|
|
if result.endswith(b' FOUND'):return 'rejected','Malware or unsafe scan condition detected'
|
|
return 'error','Scanner could not verify this file'
|
|
|
|
def scan_one(storage, scanner=None):
|
|
scanner=scanner or ClamAV()
|
|
with connect() as c:
|
|
row=c.execute('''SELECT * FROM dtf_local.uploads WHERE complete AND purged_at IS NULL
|
|
AND expires_at>now() AND scan_state IN ('pending','error') AND scan_after<=now()
|
|
ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1''').fetchone()
|
|
if not row:return False
|
|
try:
|
|
stream=storage.client.get_object(Bucket=storage.bucket,Key=row['object_key'])['Body']
|
|
try:state,reason=scanner.scan(stream,row['size'])
|
|
finally:stream.close()
|
|
except Exception:
|
|
state,reason='error','Malware scanner unavailable; file remains blocked'
|
|
c.execute("""UPDATE dtf_local.uploads SET scan_state=%s,scan_reason=%s,scanned_at=now(),
|
|
scan_after=now()+interval '1 minute',
|
|
expires_at=CASE WHEN %s IN ('rejected','error') THEN LEAST(expires_at,now()+interval '3 days') ELSE expires_at END
|
|
WHERE id=%s""",(state,reason,state,row['id']))
|
|
audit('artwork_scan', upload=str(row['id']), result=state)
|
|
return True
|
|
|
|
def scan_loop(storage):
|
|
while True:
|
|
try:
|
|
if scan_one(storage):continue
|
|
except Exception:
|
|
audit('scanner_worker_error')
|
|
time.sleep(1)
|