From 35590e5f9a95288ba32385a92a39f032293a63fe Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 15:28:34 -0400 Subject: [PATCH 01/12] added arary of execution flags --- database/fuzzili.sql | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) create mode 100644 database/fuzzili.sql diff --git a/database/fuzzili.sql b/database/fuzzili.sql new file mode 100644 index 000000000..124eaa1d0 --- /dev/null +++ b/database/fuzzili.sql @@ -0,0 +1,32 @@ +main=> CREATE TABLE main ( + fuzzer_id SERIAL PRIMARY KEY, + created_at TIMESTAMP DEFAULT NOW() +); + +CREATE TABLE fuzzer ( + program_base64 TEXT PRIMARY KEY, + fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, + inserted_at TIMESTAMP DEFAULT NOW() +); + +CREATE TABLE program ( + program_base64 TEXT PRIMARY KEY, + fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, + created_at TIMESTAMP DEFAULT NOW() +); + +CREATE TABLE execution ( + execution_id SERIAL PRIMARY KEY, + program_base64 TEXT NOT NULL REFERENCES program(fuzzer_id) ON DELETE CASCADE, + feedback_vector JSONB, + turboshaft_ir TEXT, + coverage_total NUMERIC(5,2), + created_at TIMESTAMP DEFAULT NOW(), + execution_flags TEXT[] +); + + +ALTER TABLE program +ADD CONSTRAINT fk_program_fuzzer +FOREIGN KEY (program_base64) +REFERENCES fuzzer(program_base64) \ No newline at end of file From a9eedb1fdc6929c99f64325429ab9403a93cd4e4 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 15:42:16 -0400 Subject: [PATCH 02/12] added comments --- database/fuzzili.sql | 36 ++++++++++++++++++------------------ 1 file changed, 18 insertions(+), 18 deletions(-) diff --git a/database/fuzzili.sql b/database/fuzzili.sql index 124eaa1d0..7b3f412ef 100644 --- a/database/fuzzili.sql +++ b/database/fuzzili.sql @@ -1,32 +1,32 @@ -main=> CREATE TABLE main ( - fuzzer_id SERIAL PRIMARY KEY, - created_at TIMESTAMP DEFAULT NOW() +CREATE TABLE main ( + fuzzer_id SERIAL PRIMARY KEY, + created_at TIMESTAMP DEFAULT NOW() ); CREATE TABLE fuzzer ( - program_base64 TEXT PRIMARY KEY, - fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, - inserted_at TIMESTAMP DEFAULT NOW() + program_base64 TEXT PRIMARY KEY, -- Base64-encoded fuzzer program (unique identifier) + fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, -- Links to parent fuzzer instance + inserted_at TIMESTAMP DEFAULT NOW() ); +-- Program table: Stores generated test programs CREATE TABLE program ( - program_base64 TEXT PRIMARY KEY, - fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, - created_at TIMESTAMP DEFAULT NOW() + program_base64 TEXT PRIMARY KEY, -- Base64-encoded test program (unique identifier) + fuzzer_id INT NOT NULL REFERENCES main(fuzzer_id) ON DELETE CASCADE, -- Links to parent fuzzer instance + created_at TIMESTAMP DEFAULT NOW() ); CREATE TABLE execution ( - execution_id SERIAL PRIMARY KEY, - program_base64 TEXT NOT NULL REFERENCES program(fuzzer_id) ON DELETE CASCADE, - feedback_vector JSONB, - turboshaft_ir TEXT, - coverage_total NUMERIC(5,2), - created_at TIMESTAMP DEFAULT NOW(), - execution_flags TEXT[] + execution_id SERIAL PRIMARY KEY, -- Unique identifier for each execution + program_base64 TEXT NOT NULL REFERENCES program(program_base64) ON DELETE CASCADE, -- Links to the executed program + feedback_vector JSONB, -- JSON structure containing execution feedback data + turboshaft_ir TEXT, -- Turboshaft intermediate representation output + coverage_total NUMERIC(5,2), -- Total code coverage percentage (0.00 to 999.99) + created_at TIMESTAMP DEFAULT NOW(), -- Timestamp when execution occurred + execution_flags TEXT[] -- Array of flags/options used during execution ); - ALTER TABLE program ADD CONSTRAINT fk_program_fuzzer FOREIGN KEY (program_base64) -REFERENCES fuzzer(program_base64) \ No newline at end of file +REFERENCES fuzzer(program_base64); \ No newline at end of file From 99f9b310581555c8a4ecbbea28b8b1822a3d24d9 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 15:55:20 -0400 Subject: [PATCH 03/12] added execution_type states --- database/fuzzili.sql | 21 ++++++++++++++++++++- 1 file changed, 20 insertions(+), 1 deletion(-) diff --git a/database/fuzzili.sql b/database/fuzzili.sql index 7b3f412ef..12b06a3f6 100644 --- a/database/fuzzili.sql +++ b/database/fuzzili.sql @@ -9,6 +9,14 @@ CREATE TABLE fuzzer ( inserted_at TIMESTAMP DEFAULT NOW() ); +CREATE TABLE execution_type ( + id SERIAL PRIMARY KEY, + title VARCHAR(32) NOT NULL +); + +--preseed with AI Mut, Delta, OG Testcase type +-- reference in execution table + -- Program table: Stores generated test programs CREATE TABLE program ( program_base64 TEXT PRIMARY KEY, -- Base64-encoded test program (unique identifier) @@ -19,6 +27,7 @@ CREATE TABLE program ( CREATE TABLE execution ( execution_id SERIAL PRIMARY KEY, -- Unique identifier for each execution program_base64 TEXT NOT NULL REFERENCES program(program_base64) ON DELETE CASCADE, -- Links to the executed program + execution_type_id INTEGER NOT NULL REFERENCES execution_type(id), -- Links to execution type feedback_vector JSONB, -- JSON structure containing execution feedback data turboshaft_ir TEXT, -- Turboshaft intermediate representation output coverage_total NUMERIC(5,2), -- Total code coverage percentage (0.00 to 999.99) @@ -29,4 +38,14 @@ CREATE TABLE execution ( ALTER TABLE program ADD CONSTRAINT fk_program_fuzzer FOREIGN KEY (program_base64) -REFERENCES fuzzer(program_base64); \ No newline at end of file +REFERENCES fuzzer(program_base64); + +INSERT INTO execution_type (title) VALUES + ('AI Mutation'), + ('Delta Analysis'), + ('Testcase'); + +CREATE INDEX idx_execution_program ON execution(program_base64); +CREATE INDEX idx_execution_type ON execution(execution_type_id); +CREATE INDEX idx_execution_created ON execution(created_at); +CREATE INDEX idx_execution_coverage ON execution(coverage_total); \ No newline at end of file From da0557ec3f9f7c469c2a7312ec52f6e65116f5a2 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 15:58:30 -0400 Subject: [PATCH 04/12] added --- database/fuzzili.sql | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/database/fuzzili.sql b/database/fuzzili.sql index 12b06a3f6..4bfa0576d 100644 --- a/database/fuzzili.sql +++ b/database/fuzzili.sql @@ -11,7 +11,7 @@ CREATE TABLE fuzzer ( CREATE TABLE execution_type ( id SERIAL PRIMARY KEY, - title VARCHAR(32) NOT NULL + title VARCHAR(32) NOT NULL UNIQUE ); --preseed with AI Mut, Delta, OG Testcase type @@ -41,9 +41,10 @@ FOREIGN KEY (program_base64) REFERENCES fuzzer(program_base64); INSERT INTO execution_type (title) VALUES - ('AI Mutation'), - ('Delta Analysis'), - ('Testcase'); + ('ai_mutation'), + ('delta_analysis'), + ('directed_testcases'), + ('generalistic_testcases'); CREATE INDEX idx_execution_program ON execution(program_base64); CREATE INDEX idx_execution_type ON execution(execution_type_id); From 52d9699cf2c39474124efe6b3bd8611b456919d9 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 15:58:45 -0400 Subject: [PATCH 05/12] a --- database/fuzzili.sql | 3 --- 1 file changed, 3 deletions(-) diff --git a/database/fuzzili.sql b/database/fuzzili.sql index 4bfa0576d..27af15141 100644 --- a/database/fuzzili.sql +++ b/database/fuzzili.sql @@ -14,9 +14,6 @@ CREATE TABLE execution_type ( title VARCHAR(32) NOT NULL UNIQUE ); ---preseed with AI Mut, Delta, OG Testcase type --- reference in execution table - -- Program table: Stores generated test programs CREATE TABLE program ( program_base64 TEXT PRIMARY KEY, -- Base64-encoded test program (unique identifier) From 2ed37141866136068a271a6a9cd884e0f734b38d Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 16:34:55 -0400 Subject: [PATCH 06/12] added comments --- database/fuzzili.sql | 2 ++ 1 file changed, 2 insertions(+) diff --git a/database/fuzzili.sql b/database/fuzzili.sql index 27af15141..508c4f381 100644 --- a/database/fuzzili.sql +++ b/database/fuzzili.sql @@ -43,6 +43,8 @@ INSERT INTO execution_type (title) VALUES ('directed_testcases'), ('generalistic_testcases'); + +-- Indexes for performance, just query indexes to save on speed CREATE INDEX idx_execution_program ON execution(program_base64); CREATE INDEX idx_execution_type ON execution(execution_type_id); CREATE INDEX idx_execution_created ON execution(created_at); From a62ef94dd658492371bf80685e9533611cbf6d98 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 18:17:59 -0400 Subject: [PATCH 07/12] added new sql --- database/{fuzzili.sql => database.sql} | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) rename database/{fuzzili.sql => database.sql} (96%) diff --git a/database/fuzzili.sql b/database/database.sql similarity index 96% rename from database/fuzzili.sql rename to database/database.sql index 508c4f381..b3f7b1fb8 100644 --- a/database/fuzzili.sql +++ b/database/database.sql @@ -10,8 +10,8 @@ CREATE TABLE fuzzer ( ); CREATE TABLE execution_type ( - id SERIAL PRIMARY KEY, - title VARCHAR(32) NOT NULL UNIQUE + id SERIAL PRIMARY KEY, + title VARCHAR(32) NOT NULL UNIQUE ); -- Program table: Stores generated test programs From f1bd78ef3000498944bd73461e96e88408cf7e84 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 19:01:00 -0400 Subject: [PATCH 08/12] added multithreaded database insert --- vrig_docker/requirements.txt | 4 + vrig_docker/sync.py | 271 +++++++++++++++++++++++++++++++---- 2 files changed, 248 insertions(+), 27 deletions(-) create mode 100644 vrig_docker/requirements.txt diff --git a/vrig_docker/requirements.txt b/vrig_docker/requirements.txt new file mode 100644 index 000000000..eeca2902f --- /dev/null +++ b/vrig_docker/requirements.txt @@ -0,0 +1,4 @@ +asyncpg==0.29.0 +redis[hiredis]==5.0.1 +asyncio-mqtt==0.16.1 +backoff==2.2.1 diff --git a/vrig_docker/sync.py b/vrig_docker/sync.py index cddfd5215..b1032e5f4 100644 --- a/vrig_docker/sync.py +++ b/vrig_docker/sync.py @@ -1,30 +1,220 @@ import asyncio, os, time from redis.asyncio import Redis import asyncpg +from typing import List, Dict, Any GROUP = os.getenv("GROUP", "g_fuzz") CONSUMER = os.getenv("CONSUMER", "c_sync_1") STREAMS = os.getenv("STREAMS", "redis1=redis://redis1:6379,redis2=redis://redis2:6379").split(",") STREAM_NAME = "stream:fuzz:updates" PG_DSN = os.getenv("PG_DSN", "postgres://fuzzuser:pass@pg:5432/main") +DB_WORKER_THREADS = int(os.getenv("DB_WORKER_THREADS", "4")) +BATCH_SIZE = int(os.getenv("BATCH_SIZE", "100")) +BATCH_TIMEOUT = float(os.getenv("BATCH_TIMEOUT", "1.0")) -CREATE_GROUP_OK = {"OK", "BUSYGROUP Consumer Group name already exists"} - -UPSERT_SQL = """ -INSERT INTO program (program_base64, fuzzer_id, feedback_vector, turboshaft_ir, coverage_total, created_at) -VALUES ($1, $2, $3, $4, $5, NOW()) +UPSERT_PROGRAM_SQL = """ +INSERT INTO program (program_base64, fuzzer_id, created_at) +VALUES ($1, $2, NOW()) ON CONFLICT (program_base64) DO UPDATE SET + created_at = NOW(); +""" + +UPSERT_EXECUTION_SQL = """ +INSERT INTO execution (program_base64, execution_type_id, feedback_vector, turboshaft_ir, coverage_total, execution_flags, created_at) +VALUES ($1, $2, $3, $4, $5, $6, NOW()) +ON CONFLICT (program_base64, execution_type_id) DO UPDATE SET feedback_vector = EXCLUDED.feedback_vector, turboshaft_ir = EXCLUDED.turboshaft_ir, coverage_total = EXCLUDED.coverage_total, + execution_flags = EXCLUDED.execution_flags, created_at = NOW(); """ UPDATE_FEEDBACK_SQL = """ -UPDATE program SET feedback_vector = $2 +UPDATE execution SET feedback_vector = $2 WHERE program_base64 = $1; """ +DELETE_SQL = """ +DELETE FROM program WHERE program_base64 = $1; +""" + +class DatabaseWorker: + def __init__(self, worker_id: int, dsn: str): + self.worker_id = worker_id + self.dsn = dsn + self.connection_pool = None + + async def initialize(self): + self.connection_pool = await asyncpg.create_pool( + self.dsn, + min_size=2, + max_size=10, + command_timeout=30 + ) + + async def process_batch(self, batch: List[Dict[str, Any]]): + if not batch: + return + + async with self.connection_pool.acquire() as conn: + program_upserts = [] + execution_upserts = [] + updates = [] + deletes = [] + + for operation in batch: + op_type = operation.get('op') + if op_type == 'del': + deletes.append(operation) + elif op_type == 'update_feedback': + updates.append(operation) + else: + if operation.get('execution_type_id'): + execution_upserts.append(operation) + else: + program_upserts.append(operation) + + if program_upserts: + await self._batch_upsert_program(conn, program_upserts) + if execution_upserts: + await self._batch_upsert_execution(conn, execution_upserts) + if updates: + await self._batch_update_feedback(conn, updates) + if deletes: + await self._batch_delete(conn, deletes) + + async def _batch_upsert_program(self, conn, operations: List[Dict[str, Any]]): + if not operations: + return + + values = [] + for op in operations: + program_base64 = op.get('program_base64', '') + fuzzer_id = int(op.get('fuzzer_id', 0) or 0) + values.append((program_base64, fuzzer_id)) + + await conn.executemany(UPSERT_PROGRAM_SQL, values) + + async def _batch_upsert_execution(self, conn, operations: List[Dict[str, Any]]): + if not operations: + return + + values = [] + for op in operations: + program_base64 = op.get('program_base64', '') + execution_type_id = int(op.get('execution_type_id', 1) or 1) + feedback_vector = op.get('feedback_vector', 'null') + turboshaft_ir = op.get('turboshaft_ir', '') + coverage_total = float(op.get('coverage_total', 0) or 0) + execution_flags = op.get('execution_flags', []) + + values.append((program_base64, execution_type_id, feedback_vector, turboshaft_ir, coverage_total, execution_flags)) + + await conn.executemany(UPSERT_EXECUTION_SQL, values) + + async def _batch_update_feedback(self, conn, operations: List[Dict[str, Any]]): + if not operations: + return + + values = [] + for op in operations: + program_base64 = op.get('program_base64', '') + feedback_vector = op.get('feedback_vector', 'null') + values.append((program_base64, feedback_vector)) + + await conn.executemany(UPDATE_FEEDBACK_SQL, values) + + async def _batch_delete(self, conn, operations: List[Dict[str, Any]]): + if not operations: + return + + values = [(op.get('program_base64', ''),) for op in operations] + await conn.executemany(DELETE_SQL, values) + + async def close(self): + if self.connection_pool: + await self.connection_pool.close() + +class DatabaseBatchProcessor: + def __init__(self, dsn: str, num_workers: int = DB_WORKER_THREADS): + self.dsn = dsn + self.num_workers = num_workers + self.workers: List[DatabaseWorker] = [] + self.operation_queue = asyncio.Queue() + self.running = False + + async def initialize(self): + self.workers = [] + for i in range(self.num_workers): + worker = DatabaseWorker(i, self.dsn) + await worker.initialize() + self.workers.append(worker) + + self.running = True + + def add_operation(self, operation: Dict[str, Any]): + self.operation_queue.put_nowait(operation) + + async def process_operations(self): + batch = [] + last_batch_time = time.time() + + while self.running: + try: + try: + operation = await asyncio.wait_for( + self.operation_queue.get(), + timeout=0.1 + ) + batch.append(operation) + except asyncio.TimeoutError: + pass + + current_time = time.time() + should_process_batch = ( + len(batch) >= BATCH_SIZE or + (batch and current_time - last_batch_time >= BATCH_TIMEOUT) + ) + + if should_process_batch and batch: + worker_index = hash(batch[0].get('program_base64', '')) % len(self.workers) + worker = self.workers[worker_index] + + try: + await worker.process_batch(batch.copy()) + except Exception as e: + pass + + batch.clear() + last_batch_time = current_time + + await asyncio.sleep(0.01) + + except Exception as e: + await asyncio.sleep(0.1) + + async def close(self): + self.running = False + + remaining_ops = [] + while not self.operation_queue.empty(): + try: + remaining_ops.append(self.operation_queue.get_nowait()) + except: + break + + if remaining_ops: + for worker in self.workers: + try: + await worker.process_batch(remaining_ops) + break + except: + pass + + for worker in self.workers: + await worker.close() + async def ensure_group(r: Redis, stream: str): try: await r.xgroup_create(stream, GROUP, id="$", mkstream=True) @@ -32,50 +222,77 @@ async def ensure_group(r: Redis, stream: str): if "BUSYGROUP" not in str(e): raise -async def consume_stream(label: str, redis_url: str, pg): +async def consume_stream(label: str, redis_url: str, batch_processor: DatabaseBatchProcessor): r = Redis.from_url(redis_url) await ensure_group(r, STREAM_NAME) + while True: try: - # Read new messages for this consumer resp = await r.xreadgroup(GROUP, CONSUMER, {STREAM_NAME: ">"}, count=100, block=5000) if not resp: continue - # resp = [(b'stream:fuzz:updates', [(id, {b'k':b'v', ...}), ...])] + for _, entries in resp: for msg_id, data in entries: op = data.get(b'op', b'').decode() program_base64 = data.get(b'program_base64', b'').decode() + operation = { + 'op': op, + 'program_base64': program_base64, + 'msg_id': msg_id + } + if op == "del": - # Delete program entry - await pg.execute( - "DELETE FROM program WHERE program_base64=$1", program_base64 - ) + pass elif op == "update_feedback": - # Update only the feedback_vector field - feedback_vector = data.get(b'feedback_vector', b'null').decode() - await pg.execute(UPDATE_FEEDBACK_SQL, program_base64, feedback_vector) + operation['feedback_vector'] = data.get(b'feedback_vector', b'null').decode() else: - # Full upsert (op == "set" or default) - fuzzer_id = int(data.get(b'fuzzer_id', b'0').decode() or 0) - feedback_vector = data.get(b'feedback_vector', b'null').decode() - turboshaft_ir = data.get(b'turboshaft_ir', b'').decode() - coverage_total = float(data.get(b'coverage_total', b'0').decode() or 0) + operation['fuzzer_id'] = data.get(b'fuzzer_id', b'0').decode() + operation['execution_type_id'] = data.get(b'execution_type_id', b'1').decode() + operation['feedback_vector'] = data.get(b'feedback_vector', b'null').decode() + operation['turboshaft_ir'] = data.get(b'turboshaft_ir', b'').decode() + operation['coverage_total'] = data.get(b'coverage_total', b'0').decode() - await pg.execute(UPSERT_SQL, program_base64, fuzzer_id, feedback_vector, turboshaft_ir, coverage_total) + execution_flags_str = data.get(b'execution_flags', b'').decode() + if execution_flags_str: + operation['execution_flags'] = execution_flags_str.split(',') + else: + operation['execution_flags'] = [ + 'is_debug=false', + 'v8_enable_i18n_support=false', + 'dcheck_always_on=true', + 'v8_static_library=true', + 'v8_enable_verify_heap=true', + 'v8_fuzzilli=true', + 'sanitizer_coverage_flags=trace-pc-guard', + 'target_cpu=x64' + ] + + batch_processor.add_operation(operation) await r.xack(STREAM_NAME, GROUP, msg_id) + except Exception as e: - # backoff on errors await asyncio.sleep(1) async def main(): - pg = await asyncpg.connect(PG_DSN) - tasks = [] + batch_processor = DatabaseBatchProcessor(PG_DSN, DB_WORKER_THREADS) + await batch_processor.initialize() + + batch_task = asyncio.create_task(batch_processor.process_operations()) + + stream_tasks = [] for pair in STREAMS: label, url = pair.split("=") - tasks.append(asyncio.create_task(consume_stream(label, url, pg))) - await asyncio.gather(*tasks) + task = asyncio.create_task(consume_stream(label, url, batch_processor)) + stream_tasks.append(task) + + try: + await asyncio.gather(batch_task, *stream_tasks) + except KeyboardInterrupt: + pass + finally: + await batch_processor.close() if __name__ == "__main__": asyncio.run(main()) From f0a54796569e50c22e00d330afa94642fa63c42c Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 19:04:44 -0400 Subject: [PATCH 09/12] renamed database --- database/database.sql | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/database/database.sql b/database/database.sql index b3f7b1fb8..a69b10443 100644 --- a/database/database.sql +++ b/database/database.sql @@ -38,7 +38,7 @@ FOREIGN KEY (program_base64) REFERENCES fuzzer(program_base64); INSERT INTO execution_type (title) VALUES - ('ai_mutation'), + ('agentic_analysis'), ('delta_analysis'), ('directed_testcases'), ('generalistic_testcases'); From 7534bded6565b4eb0357a975ac6e2b955ca24c2c Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 19:16:09 -0400 Subject: [PATCH 10/12] updated magic numbers --- vrig_docker/sync.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/vrig_docker/sync.py b/vrig_docker/sync.py index b1032e5f4..7ab2a3fed 100644 --- a/vrig_docker/sync.py +++ b/vrig_docker/sync.py @@ -9,8 +9,8 @@ STREAM_NAME = "stream:fuzz:updates" PG_DSN = os.getenv("PG_DSN", "postgres://fuzzuser:pass@pg:5432/main") DB_WORKER_THREADS = int(os.getenv("DB_WORKER_THREADS", "4")) -BATCH_SIZE = int(os.getenv("BATCH_SIZE", "100")) -BATCH_TIMEOUT = float(os.getenv("BATCH_TIMEOUT", "1.0")) +BATCH_SIZE = int(os.getenv("BATCH_SIZE", "400")) +BATCH_TIMEOUT = float(os.getenv("BATCH_TIMEOUT", "0.1")) UPSERT_PROGRAM_SQL = """ INSERT INTO program (program_base64, fuzzer_id, created_at) From 772e07fd7c60b09f43c9236cde61185a09ac209b Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Thu, 9 Oct 2025 19:48:50 -0400 Subject: [PATCH 11/12] moved database into docker --- vrig_docker/__pycache__/sync.cpython-312.pyc | Bin 0 -> 17515 bytes {database => vrig_docker}/database.sql | 0 vrig_docker/sync.py | 56 +++++++-- vrig_docker/test_sync.py | 118 +++++++++++++++++++ 4 files changed, 165 insertions(+), 9 deletions(-) create mode 100644 vrig_docker/__pycache__/sync.cpython-312.pyc rename {database => vrig_docker}/database.sql (100%) create mode 100644 vrig_docker/test_sync.py diff --git a/vrig_docker/__pycache__/sync.cpython-312.pyc b/vrig_docker/__pycache__/sync.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..62e1a3617ba9513eb88e3b5cd63dabdd1c037bde GIT binary patch literal 17515 zcmeHudvFtHo?y4STQAF!EXlIur~JeMV`Jk7#x`JV$q>L0FobwW6rq-l99wc)a=;!r zVM%Jo?6Hfv*$rkVcQ{kojoD;&oZYGmsk*JfAy7=+Rh?vok?7-Uxw~A|r0!7yrz(NX zRo(CVTHTV344JE|{p&uJzW(+1e*C`Q^Xq??o2w=uoZtKVQ~#}jApR?Q5raGhtZb7J z#ASjdSaOsY6rSWDiQi>|GI+{HsWEzxCWSm^kbyL6R6eE{RE#MHm1C+w6-melFF9fo z8&5x3nq@wQ9#9y}S+|s&RY;{Y>y}coDzVf^koh4(=%n=`0ew;c_eE4DMJ zhf;a#mejI3v82?t;h8n`Y~E8fj8Mb0ZXX8L_*9L2)--5l^9Ku9^Pq(-7__pM!9vzL zSR^C*2)6KDf-Ms33>M?GVn{2I(n{D;FU^*DONVII_AWhG#+G|a*$VhqdP|1nY!#&0 z*y?wQ!E$c}Tl0D?L2worAu;02Tu}jr^zx5W$>J)I#3VUM+BbZI();Zs&m8u$zM!4r z>HWT7h^Kpfo)AxU`zLv(@9@E+hj{t0YiMF>idXg=>_2k!fa5T)^tghPevhk>R~#8Q z?C3slWU7Tjtu%Jw$M%K>;dOiK*0%OK3^%RLYQkAh$hs+GFvNM?W7|<5+b70ZcgP#$ z<%jxQy+`_|s>cJt&@kr>LX)V#M9|A^A9n|XyT*sNw={2UY8Z3-{Cr;TZr4i(5ASmv zb`9)-0rei?$!1=?yL+H#kL$?ZK?hGYZ`s1<3JC*y4>%4U9hjoF)HhB|s`vIEaU347 z?d=~pXhW^pYEe&K&gEn4Y#ta}$jiFiA-nom_x_`fBevSA#yVU7!Ix_7cJ;x28?1_R z|K6T~jB0jU??KzqL%rPt4%-pOfZCSPjIB$k+o7I%Z&lYqr?%R0E-=iw$6P1fL2qkw zoh@T9K0i!9+i0qqvX5wWTTc%wqrr@E{l%(rrP_Grz}8!(|J!k;+shc%;;LbJ&v-o( zAzuKtYxOMaY*LmhG&%0Y+c@O)vM1f1*IjRVJ)r~)0#A@^uz6GNl!;o$msfM!^=75=bORW0yu6?XzQ*r`q-fixn74z@LefhkrssbAXT zIP9>ko(jrogm$*5dmZ~7fK1NA2M>s4Y|p6SNgaCpU{WtD{{w)_Kq?TE00u#NDOTpC zfzVJuR2Y^9q9SJ*ufnVJ(q7e&%nIeve@Mp4p@f=MWX)4TUd}psD#+8Ulc$Ef+;#GD zAWyqao(A%C>*VD^o}SghO!6jmc0I4{b%)&OB)t^iUiWgm>NI`=PVf{P^n0k(+Moww zcIbVj9-_-cm<$sxfP@g1!CQ(+5x`L4EyWxOVv@K<_S-2gABu8jh&ZT>d9uaJ!xn-6&Uk2gt)-C&(rWG=Nra^&Nn*h_Tg6-8}Pu;8xX~< zAD`rvE|<^m3%OiV1*<1lFJ^3orh^uU&Jqc0#SDExn=qQ;p-hxj&d_(2#Ve3Uc%<3D z3lb{cf&a=lMEeL5<{y&sN2HlTW-5e<5n6=PPnKk5EcFTfaVqZ}SgD`FOy84Trv8{= z(5j(%<1k6I5f{|IqE3@`rvDq%#WzUk`KA;iSTE(6&+iKQrn~^T&j=7K>)4pv&$>ds zF>hcZWT$zB8?DRuFrOpP47}WU6F7vr*GR4q9zir08;3*>UX51I=N^Tor>tvMFrBpx z)j|vTd*b`E#Qo~Jc=fhu^|pjYe{s+GJ&AmCBEKL}T>7nwsZoDNFg7(_*d&6U1T|3a z+=f>28re+$O-?KD*pu#%=aeUN9TiXpuQ&@UBCL3r43d+uMkA25HjGH?md+7bdo%G4 zoaoGTo|X+0(-iz^_%raA!(Rb^Wti+GUb##V)2gs)M5GMqCPA1K1%+vKSk20?fMEzL zN6>iJgs@61OA;40L1`*gTJhY{yTjD;%7p3XmSCP+BKsU^aoj=j!q39;=e2U-3QdG5 zt`K%cwPp$66(|EGifVY;lLYfD{mF4XUDkT?XY{AQ z_42aTlSKS1{V8!hW7c}g_4@-=1b$LJs4JQH3qmSThgtH%5YUJZ;r}P_ABX=j_;6r%X85IPC+s=zqNMTn2_ioiHtnjpeu;9O!5sA?`W@9_qM!XT%L*AUiK zxx1i&;IDxsW7>L8XG&-c4~zwIgqYB#eb~-eqHD!c=s* z|5E>a`(4xKM1FC^*0`MCn3?#%WQ&_NL`@qOR9~5z5~hl{$sRS?7wlh|niHnVxM^e5 zv~i*RD^p9Nz;=21(sZPH+uefpgvB1WY>8U7EL!hcI+6u2golNMzBpknjGH${&6^k4 z#r-jJcU<2c(Rbg|!H^-NA!=?|q!;}$vm>r|MD&h(x>B66C2HQXXj*(bX6}pY`y%?j zdpaA=Xo{Mf7MmmOFT~9I;`)6N{l0rTGtOy@nj06Zt}9~ZU2*-ch<=xtQ5!YaE|?Zh z$INYUeOpA|_FXp#ONkftAzIYO$vPr$H~IYoy#xMz@ZSgjUGQ)FK8eeIj|-=3{1E&M zHtt-7a|`()=~T=AiUMe1cfPZT`eMgMh;P=C&RqFTdlA52(xg)-|B}J5jdbS8zbwaa zC+RfGzwAQrEeq+)m*28txS4bo$ZxgKD7`%gTD}b}Sn1o09KpFHq~F%zI1k509OvV> z0LNAml`N7wH!E(Jlg=v5?Ft2g?FvZe)e8O5QC8>=Zk~-`RPx<^IZVSL%BSd_TKykbvPfhZ=qlb_V>^jiBx8HT- z#r<~RBB&&ig1FtNxcp7`=!7@OZ9*x%5CxIQD9`{*0k|f3pDJH-?j>3#QOO}x_ce%8 zhu?}0f1Ytho-hUlS)EojB87y+S2#mUX;hQZEn`XnCzq^bcNx7@C*hPr=EIh!wi1d-A#=;8nUI7O>y)NgiQT6VL=r=$-M|*E9=mcZQq{b8GFH+))Bk+* z*|E?bDQvl`X-z7iiKJS>ps&idEY`%d?Qv!MvaovL~!CZy1G8rA>jM1{u1)|hsC zT)BN&x%~-u5V~Kr@kkYT2mUMn8fZgS+Hv)tiFRbCP7EuF`_qbbXp$n$w5F9IiAAkb zI>U;eqF|B&2`EUKmWLH#Wmpwf56MNQ3DoWbOzpC$Tp{cofwlj! z3^Avl!_;U$RN;_$=Nm-Nd7UW!r_6#p?VlJO<#SgV2YKyk;~*xGyiO9>UD!eh924_a zL@{}0Xf)ssWe`qN>O44(J9>@|>^&$j*^~&sF6~4E<3s=gcUm~9V|#Y=5hhpr4o zs#~wiW2HMmL@zGI9DL{E?pSf#%z^tQ74uCWwSLeV*|`0BcdT;9a%Inr@>oeP7U{*h zvRG-y%)tjKF25tPsq=c^&UEC}(TIOKQux+g&D+U5G`gfIgJf@AJQmY-#+9AR%FYzm z@1c5f$(sz>lSAK>%K=_>auJ^6le-K+)|1=*V^8inSf9b`F>jY-1L!an!hrRwczXxN zcIga?rEof`$~w9Dym^zQEN*@q>+=P`8D-qU2f z6;h|F*F!!wQ<42k)W9f2S#gh2b!9=lr_oD9k%gU+!j8L|PQ=N(R&nyCMQ2Rg5m$CB zD?3DCu%vqfxl~4WSJO*2Il!yXLed1@hrwl{1=0d8QXa6Wd1SBT18Pt}$etq*4w08+ z7Yb+)9-%auO=Ae9Y@|TMB*9R!uVhH|9?&H`vT5pc+FnX!Y%tiJG2y`iTFOs=Hxdj; zFpmt$3=3sQJPRTu!}6FQjkUq+QyTj!Yr|ioS_dG?+_O&S@ACQFzI9zUW8!yYoBYHI#Cdf4BiP+p=Txv>XdM7U?Lrv6B6F1aF4RtX?eOyx? z)zl{`LSsv4>QhGKjV2(ARd$ZXD; z#0b?3{s)NOC6Y3VshY1$67X2)SwW9SYB{4$>Z_P-V0~a0w_c|fw_Wo`?OiJbq$T&0 zDrVbt({=X7);rXVZJ+z2?FUu}NK2a4%=YWeH|QHj?o{45zNC$I99$tFEoo6R9oKtr zm~MFPY`t-6sWjSo2xW96Ejp$P8fv~n-8phcxwIp?M8c&KEoZ@FuYy@Nvod zopy!bS~&&fc|~e~DIbh4z%wsScqhEPPBOIsUZMmAH{tjD{KG)2Fik=}G7hyNdcgz~ zL-;yXmVI3$=hok$YMUVf=T^$ARa7%Wf2Fh}HaCEiOLM=h=IZp7>0i8!5{nZB#WVDu zYVhW+A}DB{n4mlYKp-f1#o$c{847(?UzaGTf`tfV%V{#b!pL)&b(PWXL70|>Wg|HP zoEeL|@H9DEhooio4v6Cd0r~(f?4q>yn5fyYjY!r%U{`XlL;pO(x&6alydlCd$NN(Y zS%L@~04TL}Jn7@*seMOy2wX(iJ5)Mt{aKCE^=Y@-=TI!@g9tN{M1d9WSe|({6ff8i zE!Ystt9?XJUGpyc`*XKnSP>A;5z7JgFdZYedek>TV;Ka>(u$`I2TMz}!w4#43jznhri%+W~;1 zm6f|>W8~osBi1k|7=NM!AIhdk|?SQkA6qC*8DB>gF zE0N6yJWFy7Dndo16~_^YUO<1Tko_z?3;$X26Umk;odsI(s9|X;8{SiCy#JB2qC_7@d+;=5=PnASkRvytl-F_9p zQOM(W;eVe_`TM{-_$t&wR|#-{66$H?%iwXGR=%aYFd{rplNsgv)|9W!G>D93o0Wt! zcvc`{y@zoW`5@6~QhmL#wl7S575&41c?E^fcfu7sMe05SHZ7zr6~jL@{cYcWo_cpD zNPaXb!XZk`A>D~rBAVhF0)YV$nzUxeYX*e75e|+Ea6C2Y^>cbu0^v1JpK=FJVWN*^ zK;$@t8zm~o8IN}yJNZDo^}0Ep366Tb<93xmvh1j;ZqywNrIZ3Z7DUK+jc0=6y#5fr zk-?jSrvekBtZS9GkJpHIZv3qG3>QSzNN*6Cb%1*a0rL>vAl8$pVQ|X6jyUji(Ob}9 za2gOzU@dyX40FG*^s4=e{TG{NbPu$Ki=+QyG$|{~wI^kYJnQ|M&GDMnXiaOvQXIF` zMJ;uLD=%hgO4usDWvC72?+B{gJi8Y>4EcqZ_g~sSZ~e-&>02dXtGYUQW%3t4oz;A8 zDwtz_{=$dg9;|yrkokwmM1J9Ne*au}ArLL=`m!Wi)}LgEyy}N)!c=?BaldlI0<+-0 zru?<$nkLfXxWUD0Ux-!i`}dy6;bZZ`uSO5Q8tZpOD_yhT?JU1~;>wAIx@bw;tTRbd z`Cw$JikCJ;OPdy3KHL83_P_6pmF|jLc10|^gdy&`wC~RkBx;&w4}NVfxjc4hY`!;U zw#W7Mh~A#;0E@}R(LWna%3)HEzXc{)7BAQoE!YI(&a>V#S&|eypCpS3Yd87$eo+;g zPsx2t@jpREe#w)hhRECa_-i9LX~`HPU;8_KEh?64{~`GIunp&`oIA;d zCWnPsDkL3ddZ{QM!CSg_I&+AdTPhq|i7ySLqmlm7NF%slBc$IVNvDdwC8H56*FX!m zOfnRkNvM3wVsz9~w>HU8?q(9f^%Q~|Noe8LR+Upu-Rfi@gDcd3q;gk%A^P)q|KsKhk&!U^F%2hjIzsg#$b zlme`*BS8iPTTCgj490lJ#zT{W(^CysI8yRhnhIiFSCnsaIq-;ryy_nYnJkE@)oLZ> zHE2I=ELkI@VU7R4Ly2GnqLhpa!srFf154pd?i!{4D?>eK(7EP>z9g=%it4N82jbQ3 z(dzbib#Jt~H&*TZO5X?M1FlA7mK|$djWq$pt?mkt9VXX0D^JSkJV>&|3+tnW z^$XLn!mhZfD`M)pXR;<0P$a1(3W_d|UK)k-t18Fx^yRusbw3BYgDtADJ)|J@`(y!` zYrdz?e*!{u{*K2zJ#$klhm?MMF6V~!xsxdJ zpa`*gjzy6JMQL*e(%l40#+wNG8H!#*krzd2o5e7CMbW454|YQIE|F+oFw@9%KP1d8Ab^#YqFnY@H7+^eh?|_I&Wx3IVT=twKgVG3Ich?GA(2244qSQ6TeAv z?O)ta9viVKTgG_t7IS<6Hl5wKDh2Fc+8@ zcL^6OIR(s#dmBY)OW?d*K=7ZU=#L>{(X!Jieb}g3ACu#Qy& zr-7Uwg2>@Hcst9n|4a2_#Ay3%W3>8D#Ml?wwq&326DlHm^sk9|u zyejfExkCySn3uDV{kNq&x||p;#p1;0?@PUy=4nKC`|W`2xcVR{1_P zofB5EDtxgNlmfpQ-qiT5YMnlESdF;4&*~$b!{#E2;QycNw3G1#(R13-qjuJ5C66`y z$k7;|(aIkv6$0eRy}833ZBu5<;RU>S`*fA`C-k75#iVBz*RbL{i$*<@C#QZ_>rL}X!&+dK%vBar7DJFBr~C*3|<0kThQ5v~CtiXpk6#cQ^> zKqGd_2a&Y7%J+^P-*}8`eDE(pNj#Kvd)N)nAFG&nYi_A>N!KZ_dtAuVLmpZI_=v+Nq-+Sf(N_oPt1@XOc)*bD_INjs zdo~YGxH)!88*+0^csi9VTQ`*2;bhh`nwM}c=Z|h zP>mJXE#}mSfnc)$FK4};0P6*%etgsy;*iM1(`VpgD1n=gqJn<}kzLFE1%ls0(Pb1} zfQYBz>$`EDbnyz%^En^f0R%n)$t-+Xlg1QmPvwnkwgXm<6Ql!#XQH;Bsq@3-3Xl}Q zr6oup9AY5~Br2e;@?Y^%z&%W{grkJ`WPS-@}ng?oaQdY03 zOjrxA7GEj;(EG{pf17jTNTh9Vha&?&iM5@Gw|Sy%o>-fAP8F%xx@>8F(Ae^sTXU9K1 z9@)_!Idmk}G7xY1NwnoBv6fe$hlVJaJX*JZcIwkpk)4OY5D{yADc<^OwDr|k>uaco zrl@7>gZidlpSpGm*7eTTSpC6x{jq5Mu~_}_h-K5gqWVNff4pNL+A$D0dOX(g^4#7? zP1|zuwnwz8vlPEXi^0Q3wruu zBVj)tQm{o2bBP=ss6Whx$6((uog^P82S_5n>+#n*3sPb#?;EU_H(%O3PsR)taZSbV zG!@7Qp{o1=UlTiL->U;vb&rw$gaPR9HtZ=P;3DWAFh>;15o{wNd8yo41uu6Pa*tJc zN3KS&o`#G&W)i_r0pcRnUJZ4p%Cc8M-RUIv?3CZ>VgUZF47L1Q1p~Ri)ldKK)*91|%k~om0B!?t$-sS6XYsq*F4z0b2 zzh?AedUB<*jOU*-E>89C7xbata^P49K37|2Nw*^wHXFv zOC{R1=e3SFIhMx-GqG4d<&t|sb^#S0hLLe-!vqZ&JK_+PmcDTdxr46@T5k?2^&yJT z90VhWeApd21>ZX6?ek8a3;koN3Zel#dy#|HAhl_|iuNG7u)E2K1sILpX1bKSg z?c*>n;3=dW(Ks35V~Tg&eqfmmgKOm`4UDot@^QB+wp*Lzo44o*O(BB^GW zn)}-Pi~jTexxt9GZJ}-%-wcl^+klwpO5@tfsJ1eut&S_Jqsr=J8IYOpfkeTF5%+VD zjg!f}r_=yBVRDi6Lsj#5CXiWYSy}xTBT4e(9M29`E=RD2f^>8T1b^t`vR4VN09hc&O~OccuwB77`1SaUNM8WU3|uwJki&P-V9{E0 z4Fo>yhYyMTZwi_{mKQK4VVwE^>IpdLKoBAdf*;{!1R-Yv2gHHk4xs2o6b+*21d3io z5xy^iF_}fL$OItZC_)%)>ES837O=)-^F-MxTx;3E{R5N(ydr`?^aK$$lqyQn&^n%ywBEvBqG%OrA4GffxU&$rJG&KJgV>dqtDm*b8ZKKeS?4+z@?&`ovD__Z)nGN7_068V?7QTf z*DdD5cQ-NpwzIiOnlaL6bxAFgPeWyRn&D|kYAFjmO-VIX08e95Ez{Cx)!!+nNwVP~ zSeMC$@7|&gkYxEo0>#N!WejPUs|K4DX;?6V&5DFJwvv_eRFZ(lqVR|~PpqKVM 0: + sample = await conn.fetchrow(""" + SELECT program_base64, execution_type_id, coverage_total, execution_flags + FROM execution + ORDER BY created_at DESC + LIMIT 1 + """) + + print(f"✓ Sample execution:") + print(f" Program: {sample['program_base64']}") + print(f" Type: {sample['execution_type_id']}") + print(f" Coverage: {sample['coverage_total']}%") + print(f" Flags: {sample['execution_flags']}") + + finally: + await conn.close() + +async def main(): + print("Testing multi-threaded database sync with mock data...") + + await setup_test_database() + await test_database_worker() + await test_batch_processor() + await verify_results() + + print("All tests completed") + +if __name__ == "__main__": + asyncio.run(main()) From baf981663298fd9aa616883a9289d6380d589615 Mon Sep 17 00:00:00 2001 From: Oleg Lazari Date: Fri, 10 Oct 2025 00:25:48 -0400 Subject: [PATCH 12/12] added database support wip --- Sources/Fuzzilli/Corpus/RedisCorpus.swift | 12 +- vrig_docker/Dockerfile.fuzzillai | 112 ++++++++++++++++++ vrig_docker/Dockerfile.sync | 26 +++++ vrig_docker/__pycache__/sync.cpython-312.pyc | Bin 17515 -> 17515 bytes vrig_docker/docker-compose.yml | 114 +++++++++++++++++++ vrig_docker/sync.py | 41 ++++++- 6 files changed, 301 insertions(+), 4 deletions(-) create mode 100644 vrig_docker/Dockerfile.fuzzillai create mode 100644 vrig_docker/Dockerfile.sync create mode 100644 vrig_docker/docker-compose.yml diff --git a/Sources/Fuzzilli/Corpus/RedisCorpus.swift b/Sources/Fuzzilli/Corpus/RedisCorpus.swift index d10c07fcc..ee24391a3 100644 --- a/Sources/Fuzzilli/Corpus/RedisCorpus.swift +++ b/Sources/Fuzzilli/Corpus/RedisCorpus.swift @@ -37,7 +37,17 @@ public class RedisCorpus: ComponentBase, Collection, Corpus { } eventLoopGroup = MultiThreadedEventLoopGroup(numberOfThreads: 1) if let group = eventLoopGroup { - let address = try? SocketAddress(ipAddress: "127.0.0.1", port: 6379) + // Use Redis container hostname from environment or default to localhost + let redisHost = ProcessInfo.processInfo.environment["REDIS_HOST"] ?? "127.0.0.1" + let redisPort = Int(ProcessInfo.processInfo.environment["REDIS_PORT"] ?? "6379") ?? 6379 + let dockerNetwork = ProcessInfo.processInfo.environment["DOCKER_NETWORK"] ?? "false" + + logger.info("RedisCorpus: Connecting to Redis at \(redisHost):\(redisPort)") + if dockerNetwork.lowercased() == "true" { + logger.info("RedisCorpus: Running in Docker network mode") + } + + let address = try? SocketAddress(ipAddress: redisHost, port: redisPort) let config = RedisConnectionPool.Configuration( initialServerConnectionAddresses: [address!], maximumConnectionCount: .maximumActiveConnections(1), diff --git a/vrig_docker/Dockerfile.fuzzillai b/vrig_docker/Dockerfile.fuzzillai new file mode 100644 index 000000000..105e5b5a9 --- /dev/null +++ b/vrig_docker/Dockerfile.fuzzillai @@ -0,0 +1,112 @@ +FROM ubuntu:24.04 + +ENV DEBIAN_FRONTEND=noninteractive +ENV SWIFT_VERSION=6.2 + +# Install system dependencies including V8 build requirements +RUN apt-get update && apt-get install -y \ + git \ + wget \ + curl \ + build-essential \ + cmake \ + ninja-build \ + clang \ + llvm \ + python3 \ + python3-pip \ + python3-venv \ + pkg-config \ + unzip \ + redis-tools \ + libnss3-dev \ + libatk-bridge2.0-dev \ + libdrm2 \ + libxcomposite1 \ + libxdamage1 \ + libxrandr2 \ + libgbm1 \ + libxss1 \ + libasound2t64 \ + && rm -rf /var/lib/apt/lists/* + +# Install Swift 6.2 (using Ubuntu 22.04 package as 24.04 is not available yet) +RUN wget -q https://download.swift.org/swift-${SWIFT_VERSION}-release/ubuntu2204/swift-${SWIFT_VERSION}-RELEASE/swift-${SWIFT_VERSION}-RELEASE-ubuntu22.04.tar.gz \ + && tar xzf swift-${SWIFT_VERSION}-RELEASE-ubuntu22.04.tar.gz \ + && mv swift-${SWIFT_VERSION}-RELEASE-ubuntu22.04 /opt/swift \ + && rm swift-${SWIFT_VERSION}-RELEASE-ubuntu22.04.tar.gz + +ENV PATH="/opt/swift/usr/bin:${PATH}" + +WORKDIR /app + +# Clone and setup FuzzilliAI +RUN git clone https://github.com/VRIG-Ritsec/fuzzillai.git . + +# Build FuzzilliAI +RUN swift build -c release + +# Create necessary directories and clean up any old corpus +RUN mkdir -p ./Corpus ./logs +RUN rm -rf ./Corpus/old_corpus ./Corpus/corpus + +# Create a startup script that runs FuzzilliAI with Redis corpus +RUN echo '#!/bin/bash\n\ +# Start FuzzilliAI with Redis corpus support\n\ +\n\ +# Check if we are in Docker network\n\ +if [ "${DOCKER_NETWORK}" = "true" ]; then\n\ + echo "Running in Docker network mode"\n\ + echo "Service: ${SERVICE_NAME:-fuzzillai}"\n\ + echo "Redis Host: ${REDIS_HOST:-redis}"\n\ + echo "Redis Port: ${REDIS_PORT:-6379}"\n\ + echo "Profile: ${FUZZILLI_PROFILE:-v8}"\n\ + echo "Engine: ${FUZZILLI_ENGINE:-multi}"\n\ +fi\n\ +\n\ +# Check if V8 d8 binary exists in mounted directory\n\ +if [ -f "/v8/out/fuzzbuild/d8" ]; then\n\ + echo "Using mounted V8 binary from /v8/out/fuzzbuild/d8"\n\ + V8_BINARY="/v8/out/fuzzbuild/d8"\n\ +elif [ -f "/v8/out/d8" ]; then\n\ + echo "Using mounted V8 binary from /v8/out/d8"\n\ + V8_BINARY="/v8/out/d8"\n\ +else\n\ + echo "V8 binary not found in mounted directory. Available files:"\n\ + find /v8 -name "d8" -type f 2>/dev/null || echo "No d8 binary found"\n\ + exit 1\n\ +fi\n\ +\n\ +# Wait for Redis to be ready\n\ +echo "Waiting for Redis to be ready..."\n\ +REDIS_HOST=${REDIS_HOST:-redis}\n\ +REDIS_PORT=${REDIS_PORT:-6379}\n\ +\n\ +for i in {1..30}; do\n\ + if redis-cli -h "$REDIS_HOST" -p "$REDIS_PORT" ping > /dev/null 2>&1; then\n\ + echo "Redis is ready at $REDIS_HOST:$REDIS_PORT!"\n\ + break\n\ + fi\n\ + echo "Redis not ready at $REDIS_HOST:$REDIS_PORT, waiting... (attempt $i/30)"\n\ + sleep 2\n\ + if [ $i -eq 30 ]; then\n\ + echo "Redis connection timeout, proceeding anyway..."\n\ + fi\n\ +done\n\ +\n\ +# Clean up any problematic corpus directories\n\ +echo "Cleaning up corpus directories..."\n\ +rm -rf ./Corpus/old_corpus\n\ +\n\ +echo "Starting FuzzilliAI..."\n\ +swift run -c release FuzzilliCli \\\n\ + --profile=${FUZZILLI_PROFILE:-v8} \\\n\ + --engine=${FUZZILLI_ENGINE:-multi} \\\n\ + --corpus=redis \\\n\ + --storagePath=./Corpus \\\n\ + --resume \\\n\ + "$V8_BINARY"' > /app/start.sh && chmod +x /app/start.sh + +EXPOSE 6379 + +CMD ["/app/start.sh"] diff --git a/vrig_docker/Dockerfile.sync b/vrig_docker/Dockerfile.sync new file mode 100644 index 000000000..a68ce831b --- /dev/null +++ b/vrig_docker/Dockerfile.sync @@ -0,0 +1,26 @@ +FROM python:3.11-slim + +# Install system dependencies +RUN apt-get update && apt-get install -y \ + gcc \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /app + +# Copy requirements and install Python dependencies +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +# Copy the sync script +COPY sync.py . + +# Create a startup script +RUN echo '#!/bin/bash\n\ +echo "Starting sync service..."\n\ +echo "PostgreSQL DSN: ${PG_DSN}"\n\ +echo "Redis Streams: ${STREAMS}"\n\ +echo "Group: ${GROUP}"\n\ +echo "Consumer: ${CONSUMER}"\n\ +python sync.py' > /app/start.sh && chmod +x /app/start.sh + +CMD ["/app/start.sh"] diff --git a/vrig_docker/__pycache__/sync.cpython-312.pyc b/vrig_docker/__pycache__/sync.cpython-312.pyc index 62e1a3617ba9513eb88e3b5cd63dabdd1c037bde..93e4993eb51cd7a447578369c613fc418d74b504 100644 GIT binary patch delta 21 bcmaFe!T7p^k?S-sFBbz4FlTM#N^}7LOB@BS delta 21 bcmaFe!T7p^k?S-sFBbz4bb4;&N^}7LO`8T3 diff --git a/vrig_docker/docker-compose.yml b/vrig_docker/docker-compose.yml new file mode 100644 index 000000000..23684d00a --- /dev/null +++ b/vrig_docker/docker-compose.yml @@ -0,0 +1,114 @@ +services: + # PostgreSQL Database + postgres: + image: postgres:15 + environment: + POSTGRES_DB: main + POSTGRES_USER: fuzzuser + POSTGRES_PASSWORD: pass + DOCKER_NETWORK: "true" + SERVICE_NAME: "postgres" + volumes: + - postgres_data:/var/lib/postgresql/data + - ./database.sql:/docker-entrypoint-initdb.d/01-init.sql + ports: + - "5432:5432" + networks: + - fuzzillai_network + healthcheck: + test: ["CMD-SHELL", "pg_isready -U fuzzuser -d main"] + interval: 10s + timeout: 5s + retries: 5 + + # Redis for FuzzilliAI communication + redis: + image: redis:7-alpine + environment: + DOCKER_NETWORK: "true" + SERVICE_NAME: "redis" + ports: + - "6379:6379" + volumes: + - redis_data:/data + networks: + - fuzzillai_network + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 10s + timeout: 5s + retries: 5 + + # FuzzilliAI Fuzzer Container + fuzzillai: + build: + context: .. + dockerfile: vrig_docker/Dockerfile.fuzzillai + environment: + - REDIS_HOST=redis + - REDIS_PORT=6379 + - FUZZILLI_PROFILE=v8 + - FUZZILLI_ENGINE=multi + - DOCKER_NETWORK=true + - SERVICE_NAME=fuzzillai + volumes: + - fuzzillai_corpus:/app/Corpus + - fuzzillai_logs:/app/logs + - /home/tropic/vrig/v8/v8:/v8:ro + - fuzzillai_corpus_data:/app/Corpus/corpus + depends_on: + redis: + condition: service_healthy + restart: unless-stopped + networks: + - fuzzillai_network + healthcheck: + test: ["CMD", "pgrep", "-f", "FuzzilliCli"] + interval: 30s + timeout: 10s + retries: 3 + start_period: 60s + + # Sync Service - Processes Redis streams and sends to PostgreSQL + sync: + build: + context: . + dockerfile: Dockerfile.sync + environment: + - PG_DSN=postgres://fuzzuser:pass@postgres:5432/main + - STREAMS=redis1=redis://redis:6379 + - GROUP=g_fuzz + - CONSUMER=c_sync_1 + - DB_WORKER_THREADS=4 + - BATCH_SIZE=400 + - BATCH_TIMEOUT=0.1 + - DOCKER_NETWORK=true + - SERVICE_NAME=sync + depends_on: + postgres: + condition: service_healthy + redis: + condition: service_healthy + restart: unless-stopped + networks: + - fuzzillai_network + healthcheck: + test: ["CMD", "pgrep", "-f", "sync.py"] + interval: 30s + timeout: 10s + retries: 3 + start_period: 60s + +volumes: + postgres_data: + redis_data: + fuzzillai_corpus: + fuzzillai_corpus_data: + fuzzillai_logs: + +networks: + fuzzillai_network: + driver: bridge + ipam: + config: + - subnet: 172.20.0.0/16 diff --git a/vrig_docker/sync.py b/vrig_docker/sync.py index f0b747252..8bd5d0da4 100644 --- a/vrig_docker/sync.py +++ b/vrig_docker/sync.py @@ -5,7 +5,7 @@ GROUP = os.getenv("GROUP", "g_fuzz") CONSUMER = os.getenv("CONSUMER", "c_sync_1") -STREAMS = os.getenv("STREAMS", "redis1=redis://redis1:6379,redis2=redis://redis2:6379").split(",") +STREAMS = os.getenv("STREAMS", "redis1=redis://redis:6379,redis2=redis://redis2:6379").split(",") STREAM_NAME = "stream:fuzz:updates" PG_DSN = os.getenv("PG_DSN", "postgres://fuzzuser:pass@pg:5432/main") DB_WORKER_THREADS = int(os.getenv("DB_WORKER_THREADS", "4")) @@ -314,8 +314,43 @@ async def consume_stream(label: str, redis_url: str, batch_processor: DatabaseBa await asyncio.sleep(1) async def main(): - batch_processor = DatabaseBatchProcessor(PG_DSN, DB_WORKER_THREADS) - await batch_processor.initialize() + # Check if we're in a Docker network + docker_network = os.getenv("DOCKER_NETWORK", "false").lower() == "true" + service_name = os.getenv("SERVICE_NAME", "sync") + + print(f"Starting {service_name} service...") + if docker_network: + print("Running in Docker network mode") + print(f"PostgreSQL DSN: {PG_DSN}") + print(f"Redis Streams: {STREAMS}") + print(f"Group: {GROUP}") + print(f"Consumer: {CONSUMER}") + + # Wait for PostgreSQL to be ready with retries + print("Waiting for PostgreSQL to be ready...") + max_retries = 60 # Increased retries + retry_count = 0 + + # Initial delay to let the network settle (longer in Docker) + initial_delay = 10 if docker_network else 5 + print(f"Initial network settle delay: {initial_delay}s") + await asyncio.sleep(initial_delay) + + while retry_count < max_retries: + try: + batch_processor = DatabaseBatchProcessor(PG_DSN, DB_WORKER_THREADS) + await batch_processor.initialize() + print("Successfully connected to PostgreSQL") + break + except Exception as e: + retry_count += 1 + print(f"Failed to connect to PostgreSQL (attempt {retry_count}/{max_retries}): {e}") + if retry_count >= max_retries: + print("Max retries reached, exiting") + return + # Progressive backoff: start with 2s, increase to 5s + delay = min(2 + (retry_count * 0.1), 5) + await asyncio.sleep(delay) batch_task = asyncio.create_task(batch_processor.process_operations())