228 lines
7.7 KiB
Python
228 lines
7.7 KiB
Python
"""Data models for Mosaic intelligence extraction."""
|
|
from dataclasses import dataclass, field, asdict
|
|
from typing import Optional
|
|
import hashlib
|
|
import json
|
|
|
|
|
|
@dataclass
|
|
class Document:
|
|
id: str = "" # SHA256 of content
|
|
source: str = "" # Source profile name
|
|
source_url: str = ""
|
|
doc_type: str = "" # html, pdf, email, cable, text
|
|
title: str = ""
|
|
date: Optional[str] = None
|
|
classification: Optional[str] = None
|
|
raw_path: str = "" # Hash-sharded path to raw file
|
|
text: str = "" # Extracted plain text
|
|
metadata: dict = field(default_factory=dict)
|
|
content_hash: str = "" # SHA256 verified on every access
|
|
fetch_date: str = ""
|
|
etag: Optional[str] = None
|
|
last_modified: Optional[str] = None
|
|
char_count: int = 0
|
|
status: str = "PENDING" # PENDING|DOWNLOADING|CACHED|DOWNLOAD_FAILED|IN_PROGRESS|PARSED|PARSE_FAILED|NEEDS_OCR
|
|
parse_status: str = "PENDING" # PENDING|IN_PROGRESS|PARSED|PARSE_FAILED|NEEDS_OCR
|
|
last_extracted_at: Optional[str] = None
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
d['metadata'] = json.dumps(d['metadata'])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'Document':
|
|
row = dict(row)
|
|
if isinstance(row.get('metadata'), str):
|
|
try:
|
|
row['metadata'] = json.loads(row['metadata'])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row['metadata'] = {}
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class ExtractedTool:
|
|
name: str = ""
|
|
aliases: list = field(default_factory=list)
|
|
description: str = ""
|
|
capability: str = "" # exploit, implant, framework, utility
|
|
target_platforms: list = field(default_factory=list)
|
|
target_software: list = field(default_factory=list)
|
|
cves: list = field(default_factory=list)
|
|
source_doc_ids: list = field(default_factory=list)
|
|
source_context: str = ""
|
|
mitre_techniques: list = field(default_factory=list)
|
|
dedup_key: str = ""
|
|
|
|
def compute_dedup_key(self) -> str:
|
|
raw = f"{self.name.lower().strip()}:{self.capability.lower().strip()}"
|
|
self.dedup_key = hashlib.sha256(raw.encode()).hexdigest()[:16]
|
|
return self.dedup_key
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
for k in ('aliases', 'target_platforms', 'target_software', 'cves',
|
|
'source_doc_ids', 'mitre_techniques'):
|
|
d[k] = json.dumps(d[k])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'ExtractedTool':
|
|
row = dict(row)
|
|
for k in ('aliases', 'target_platforms', 'target_software', 'cves',
|
|
'source_doc_ids', 'mitre_techniques'):
|
|
if isinstance(row.get(k), str):
|
|
try:
|
|
row[k] = json.loads(row[k])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row[k] = []
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class TTP:
|
|
technique: str = ""
|
|
category: str = "" # initial_access, persistence, exfil, etc.
|
|
description: str = ""
|
|
actor: Optional[str] = None
|
|
mitre_id: Optional[str] = None
|
|
source_doc_ids: list = field(default_factory=list)
|
|
source_context: str = ""
|
|
dedup_key: str = ""
|
|
|
|
def compute_dedup_key(self) -> str:
|
|
raw = f"{self.technique.lower().strip()}:{self.category.lower().strip()}"
|
|
self.dedup_key = hashlib.sha256(raw.encode()).hexdigest()[:16]
|
|
return self.dedup_key
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
d['source_doc_ids'] = json.dumps(d['source_doc_ids'])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'TTP':
|
|
row = dict(row)
|
|
if isinstance(row.get('source_doc_ids'), str):
|
|
try:
|
|
row['source_doc_ids'] = json.loads(row['source_doc_ids'])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row['source_doc_ids'] = []
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class Tradecraft:
|
|
method: str = ""
|
|
domain: str = "" # humint, sigint, cyber, surveillance
|
|
description: str = ""
|
|
operational_notes: str = ""
|
|
countermeasures: str = ""
|
|
source_doc_ids: list = field(default_factory=list)
|
|
source_context: str = ""
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
d['source_doc_ids'] = json.dumps(d['source_doc_ids'])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'Tradecraft':
|
|
row = dict(row)
|
|
if isinstance(row.get('source_doc_ids'), str):
|
|
try:
|
|
row['source_doc_ids'] = json.loads(row['source_doc_ids'])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row['source_doc_ids'] = []
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class Entity:
|
|
name: str = ""
|
|
type: str = "" # person, org, program, codename, unit
|
|
description: str = ""
|
|
relationships: list = field(default_factory=list)
|
|
source_doc_ids: list = field(default_factory=list)
|
|
source_context: str = ""
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
d['relationships'] = json.dumps(d['relationships'])
|
|
d['source_doc_ids'] = json.dumps(d['source_doc_ids'])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'Entity':
|
|
row = dict(row)
|
|
for k in ('relationships', 'source_doc_ids'):
|
|
if isinstance(row.get(k), str):
|
|
try:
|
|
row[k] = json.loads(row[k])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row[k] = []
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class Infrastructure:
|
|
indicator: str = ""
|
|
indicator_type: str = "" # ipv4, ipv6, domain, url, hash
|
|
context: str = ""
|
|
source_doc_ids: list = field(default_factory=list)
|
|
source_context: str = ""
|
|
collection: str = ""
|
|
doc_title: str = ""
|
|
doc_date: Optional[str] = None
|
|
|
|
def to_dict(self) -> dict:
|
|
d = asdict(self)
|
|
d['source_doc_ids'] = json.dumps(d['source_doc_ids'])
|
|
return d
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'Infrastructure':
|
|
row = dict(row)
|
|
if isinstance(row.get('source_doc_ids'), str):
|
|
try:
|
|
row['source_doc_ids'] = json.loads(row['source_doc_ids'])
|
|
except (json.JSONDecodeError, TypeError):
|
|
row['source_doc_ids'] = []
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class ExtractionStatus:
|
|
"""Tracks per-document, per-extractor status."""
|
|
doc_id: str = ""
|
|
extractor: str = "" # tool, ttp, tradecraft, surveillance, infrastructure, entity
|
|
status: str = "PENDING" # PENDING|IN_PROGRESS|DONE|FAILED
|
|
error: Optional[str] = None
|
|
last_run: Optional[str] = None
|
|
token_count: int = 0
|
|
|
|
def to_dict(self) -> dict:
|
|
return asdict(self)
|
|
|
|
@classmethod
|
|
def from_row(cls, row: dict) -> 'ExtractionStatus':
|
|
row = dict(row)
|
|
return cls(**{k: v for k, v in row.items() if k in cls.__dataclass_fields__})
|
|
|
|
|
|
@dataclass
|
|
class TokenUsage:
|
|
"""Track API token usage per call."""
|
|
timestamp: str = ""
|
|
source: str = ""
|
|
doc_id: str = ""
|
|
extractor: str = ""
|
|
input_tokens: int = 0
|
|
output_tokens: int = 0
|
|
model: str = ""
|
|
|
|
def to_dict(self) -> dict:
|
|
return asdict(self)
|