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
17 changes: 17 additions & 0 deletions .github/workflows/push_hub.yml
Original file line number Diff line number Diff line change
Expand Up @@ -41,3 +41,20 @@ jobs:
repository: mqueryci/mquery-daemon
tags: ${{ github.sha }}
push: ${{ github.event_name == 'push' }}
build_indexer:
name: Build image
runs-on: ubuntu-latest
env:
DOCKER_BUILDKIT: 1
steps:
- name: Check out repository
uses: actions/checkout@v2
- name: Build and push the image
uses: docker/build-push-action@v1.1.0
with:
username: ${{ secrets.DOCKER_USERNAME }}
password: ${{ secrets.DOCKER_PASSWORD }}
dockerfile: ./deploy/docker/indexer.Dockerfile
repository: mqueryci/mquery-indexer
tags: ${{ github.sha }}
push: ${{ github.event_name == 'push' }}
2 changes: 1 addition & 1 deletion deploy/docker/dev.daemon.Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -9,4 +9,4 @@ RUN pip install watchdog
COPY requirements.txt src/plugins/requirements-*.txt /tmp/
RUN ls /tmp/requirements*.txt | xargs -i,, pip --no-cache-dir install -r ,,

CMD pip install -e /usr/src/app && watchmedo auto-restart --pattern=*.py --recursive -- mquery-daemon --with-indexer
CMD pip install -e /usr/src/app && watchmedo auto-restart --pattern=*.py --recursive -- mquery-daemon
12 changes: 12 additions & 0 deletions deploy/docker/dev.indexer.Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
FROM python:3.10

WORKDIR /usr/src/app/src

RUN apt update; apt install -y cmake

# mquery and plugin requirements
RUN pip install watchdog
COPY requirements.txt src/plugins/requirements-*.txt /tmp/
RUN ls /tmp/requirements*.txt | xargs -i,, pip --no-cache-dir install -r ,,

CMD pip install -e /usr/src/app && watchmedo auto-restart --pattern=*.py --recursive -- mquery-indexer
13 changes: 13 additions & 0 deletions deploy/docker/indexer.Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
FROM python:3.10

RUN apt update; apt install -y cmake

# mquery and plugin requirements
COPY requirements.txt src/plugins/requirements-*.txt /tmp/
RUN ls /tmp/requirements*.txt | xargs -i,, pip --no-cache-dir install -r ,,

COPY requirements.txt setup.py MANIFEST.in /app/
COPY src /app/src/
RUN pip install /app

ENTRYPOINT ["mquery-indexer"]
25 changes: 25 additions & 0 deletions docker-compose.dev.yml
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,31 @@ services:
volumes:
- "${SAMPLES_DIR}:/mnt/samples"
- .:/usr/src/app
depends_on:
dev-web:
condition: service_healthy
redis:
condition: service_started
ursadb:
condition: service_started
postgres:
condition: service_healthy
environment:
- "REDIS_HOST=redis"
- "MQUERY_BACKEND=tcp://ursadb:9281"
- "MQUERY_PLUGINS=${MQUERY_PLUGINS}"
- "DATABASE_URL=postgresql://postgres:password@postgres:5432/mquery"
dev-indexer:
build:
context: .
dockerfile: deploy/docker/dev.indexer.Dockerfile
links:
- redis
- ursadb
- postgres
volumes:
- "${SAMPLES_DIR}:/mnt/samples"
- .:/usr/src/app
- "s3:/mnt/s3"
depends_on:
dev-web:
Expand Down
1 change: 1 addition & 0 deletions setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
install_requires=open("requirements.txt").read().splitlines(),
scripts=[
"src/scripts/mquery-daemon",
"src/scripts/mquery-indexer",
],
classifiers=[
"Programming Language :: Python",
Expand Down
13 changes: 1 addition & 12 deletions src/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -310,7 +310,7 @@ def backend_status_datasets() -> BackendStatusDatasetsSchema:
for agent_spec in db.get_active_agents().values():
try:
ursa = UrsaDb(agent_spec.ursadb_url)
datasets.update(ursa.topology()["result"]["datasets"])
datasets.update(ursa.datasets())
except Again:
pass

Expand Down Expand Up @@ -616,17 +616,6 @@ def delete_queued_by_id(ursadb_id: str) -> StatusSchema:
return StatusSchema(status="ok")


@app.post(
"/api/queue/{ursadb_id}/index",
response_model=StatusSchema,
tags=["queue"],
dependencies=[Depends(can_manage_queues)],
)
def index_queue(ursadb_id: str) -> StatusSchema:
db.create_indexing_task(ursadb_id)
return StatusSchema(status="ok")


# Permissionless routes.
# 1. Static routes are always publicly accessible without authorisation.
# 2. /api/server is a special route always accessible for everyone.
Expand Down
16 changes: 0 additions & 16 deletions src/daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,15 +21,6 @@ def start_worker(args: argparse.Namespace, process_index: int) -> None:
w.work()


def start_indexer(args: argparse.Namespace) -> None:
setup_logging()
logging.info("Indexer [%s] running...", args.group_id)

with Connection(Redis(app_config.redis.host, app_config.redis.port)):
w = Worker([args.group_id + ":indexer"])
w.work()


def main() -> None:
"""Spawns a new agent process. Use argv if you want to use a different
group_id (it's `default` by default).
Expand All @@ -48,11 +39,6 @@ def main() -> None:
help="Specifies the number of concurrent workers to use for yara matching.",
default=1,
)
parser.add_argument(
"--with-indexer",
action="store_true",
help="If specified, indexing worker will be created (not affected by scale).",
)

args = parser.parse_args()

Expand All @@ -63,8 +49,6 @@ def main() -> None:
children = [
Process(target=start_worker, args=(args, i)) for i in range(args.scale)
]
if args.with_indexer:
children.append(Process(target=start_indexer, args=(args,)))

for child in children:
child.start()
Expand Down
20 changes: 6 additions & 14 deletions src/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -347,14 +347,6 @@ def create_search_task(
self.__schedule(agent, tasks.start_search, job)
return job

def create_indexing_task(self, ursadb_id: str) -> None:
"""Asks indexer for `ursadb_id` to index pending files."""
from . import tasks

Queue(f"{ursadb_id}:indexer", connection=self.redis).enqueue(
tasks.index_pending_files
)

def get_job_matches(
self, job_id: JobId, offset: int = 0, limit: Optional[int] = None
) -> MatchesSchema:
Expand Down Expand Up @@ -531,9 +523,9 @@ def delete_queued_files(self, ursadb_id: str) -> None:
session.query(QueuedFile).filter_by(ursadb_id=ursadb_id).delete()
session.commit()

def get_indexing_batch(self, ursadb_id: str) -> List[QueuedFile]:
BATCH_SIZE = 100

def get_pending_files(
self, ursadb_id: str, limit: int
) -> List[QueuedFile]:
# We can only batch files with the same tags and index types, so first find
# the biggest group of files waiting to be indexed:
group_subquery = (
Expand All @@ -548,7 +540,7 @@ def get_indexing_batch(self, ursadb_id: str) -> List[QueuedFile]:
.subquery()
)

# Then return up to BATCH_SIZE files from that group.
# Then return files from that group.
with self.session() as session:
group = session.execute(select(group_subquery)).first()
if not group:
Expand All @@ -561,12 +553,12 @@ def get_indexing_batch(self, ursadb_id: str) -> List[QueuedFile]:
QueuedFile.index_types == group.index_types,
QueuedFile.tags == group.tags,
)
.limit(BATCH_SIZE)
.limit(limit)
).all()

return files

def complete_indexing_batch(self, files: list[QueuedFile]) -> None:
def remove_from_pending(self, files: list[QueuedFile]) -> None:
ids = [f.id for f in files]
with self.session() as session:
session.execute(
Expand Down
2 changes: 1 addition & 1 deletion src/e2etests/test_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,7 @@ def test_query_with_taints(add_files_to_index):
db = UrsaDb("tcp://ursadb:9281")

random_taint = os.urandom(8).hex()
for dataset_id in db.topology()["result"]["datasets"].keys():
for dataset_id in db.datasets().keys():
out = db.execute_command(
f'dataset "{dataset_id}" taint "{random_taint}";'
)
Expand Down
143 changes: 143 additions & 0 deletions src/indexer.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
#!/usr/bin/env python
import argparse
import logging
from multiprocessing import Pool
from time import sleep

from .db import Database
from .plugins import PluginManager
from .util import setup_logging
from .models.queuedfile import QueuedFile
from .lib.ursadb import UrsaDb
from .config import app_config


COMPACT_THRESHOLD = 0
"""Global variable dependeng on number of workers"""


def index_batch(
job: tuple[list[QueuedFile], list[str], list[str]],
) -> list[QueuedFile]:
batch, index_types, tags = job
assert (
batch
), "This function shouldn't be called without files, likely a bug."

paths = [f.path for f in batch]

db = Database(app_config.redis.host, app_config.redis.port)
ursa = UrsaDb(app_config.mquery.backend)
plugins = PluginManager(app_config.mquery.plugins, db)

current_datasets = len(ursa.datasets())
if current_datasets > COMPACT_THRESHOLD:
ursa.execute_command("compact smart;")

ursadb_batch = []
for file_path in paths:
final_path = plugins.filter(file_path)
if final_path is None:
logging.debug("Filtering out file %s", file_path)
continue
ursadb_batch.append(final_path)

logging.debug("Batch preprocessed, asking ursadb to index it.")

ursa.index(ursadb_batch, index_types=index_types, tags=tags)
logging.debug("Ursadb indexing completed")

plugins.cleanup()
logging.debug("Cleanup completed")

return batch


def indexer_main(group_id: str, scale: int) -> None:
"""Do the indexing in an infinite loop."""
logging.info("Starting indexer for group %s", group_id)
logging.info("Workers: %s", scale)

global COMPACT_THRESHOLD
COMPACT_THRESHOLD = scale * 20 + 40
logging.info("Compact threshold: %s", COMPACT_THRESHOLD)

db = Database(app_config.redis.host, app_config.redis.port)

# How many files should one worker index at once
BATCH_SIZE = 1000

# Don't get all the files from the database at once, to avoid huge queries
LIMIT = scale * 100 * BATCH_SIZE

while True:
pending = db.get_pending_files(group_id, LIMIT)
if not pending:
logging.debug("No pending files left.")
sleep(15) # Wait a bit to collect some files.
continue

index_types, tags = pending[0].index_types, pending[0].tags
logging.debug("We are indexing type=%s, tags=%s", index_types, tags)

batches = []
next_batch = []
for f in pending:
next_batch.append(f)
if len(next_batch) > BATCH_SIZE:
batches.append(next_batch)
next_batch = []

if next_batch:
batches.append(next_batch)

types_str = ",".join(type for type in index_types)
tags_str = ",".join(tag for tag in tags)
logging.info(
"[0/%s] Collected batches (%s files), with [%s], with tags [%s]",
len(batches),
len(pending),
types_str,
tags_str,
)

jobs = [(batch, index_types, tags) for batch in batches]

done = 0
pool = Pool(processes=scale)
for completed_batch in pool.imap_unordered(
index_batch, jobs, chunksize=1
):
done += 1
db.remove_from_pending(completed_batch)
logging.info("[%s/%s] Batch completed.", done, len(batches))

logging.info("[%s/%s] Indexing completed", done, len(batches))


def main() -> None:
"""Spawns a new indexer process. Indexer will work in the background all the time."""

parser = argparse.ArgumentParser(
description="Start mquery indexer worker."
)
parser.add_argument(
"group_id",
help="Name of the agent group to join to",
nargs="?",
default="default",
)
parser.add_argument(
"--scale",
type=int,
help="Specifies the number of concurrent workers to use for yara matching.",
default=1,
)
args = parser.parse_args()

setup_logging(logging.DEBUG)
indexer_main(args.group_id, args.scale)


if __name__ == "__main__":
main()
12 changes: 10 additions & 2 deletions src/lib/ursadb.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ def __init__(self, backend: str) -> None:
self.backend = backend

def __execute(self, command: str, recv_timeout: int = 2000) -> Json:
logging.info("Ursadb command: %s", command)
logging.debug("Ursadb command: %s", command)
context = zmq.Context()
try:
socket = context.socket(zmq.REQ)
Expand Down Expand Up @@ -105,7 +105,15 @@ def status(self) -> Json:
return self.__execute("status;")

def topology(self) -> Json:
return self.__execute("topology;")
result = self.__execute("topology;")

if "error" in result:
raise RuntimeError(result["error"])

return result["result"]

def datasets(self) -> Json:
return self.topology()["datasets"]

def execute_command(self, command: str) -> Json:
return self.__execute(command, -1)
Expand Down
2 changes: 1 addition & 1 deletion src/plugins/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ def __init__(self, spec: str, db: Database) -> None:
plugin_config = db.get_plugin_config(plugin_name)
try:
active_plugins.append(plugin_class(db, plugin_config))
logging.info("Loaded plugin %s", plugin_name)
logging.debug("Loaded plugin %s", plugin_name)
except Exception:
logging.exception("Failed to load %s plugin", plugin_name)
self.active_plugins = active_plugins
Expand Down
Loading