DnDn의 Worker 는 AWS 계정에서 변경 이력과 리소스 상태를 수집하고, 이를 표준 JSON 결과물로 정규화하여 S3에 저장하는 백엔드 실행 엔진입니다.
현재 구조에서 Worker는 단순 라이브러리가 아니라, 신호(payload)를 받으면 최종 결과 JSON을 생산하는 독립 실행 서비스를 목표로 합니다.
쉽게 말하면 Worker는 다음 역할을 합니다.
- 입력:
contracts/payload/*.json형태의 작업 요청(payload) - 수집: CloudTrail, AWS Config, 이벤트 트리거(AWS Health / Security Hub 등) 관련 정보
- 정규화:
canonical.json(WEEKLY),event.json(EVENT) - 저장:
raw/,normalized/결과물을 S3에 업로드 - 확장: AWS Health 기반 이벤트 보강, 주간 운영 점검(advisor checks), Config before/after 보강
Worker는 DnDn에서 아래 흐름을 담당합니다.
- API가 Worker payload를 생성하고 전달한다 스케줄러/EventBridge 계열 트리거는 현재 구현 기준으로 API의 job 생성 경로를 통해 Worker 실행으로 연결된다
- Worker가 대상 AWS 계정(필요시 AssumeRole)으로 접근
- CloudTrail / Config / 이벤트 trigger 정보를 수집
- 결과를 표준 포맷(JSON)으로 정리
raw/,normalized/산출물을 S3에 저장- 이후 B 파트(Report)가 이 JSON과 raw evidence를 사용해 보고서, 계획서, evidence zip을 생성
즉, Worker는 “수집 → 정규화 → 저장” 의 중심입니다. 운영에서는 SQS consumer로 계속 실행되고, 로컬에서는 payload 파일 기반으로 동일 로직을 실행합니다.
- WEEKLY
- 일정 기간(보통 지난 1주)의 변경 이력 수집
- 결과:
canonical.json
- EVENT
- 특정 이벤트 시점(
event_time ± window_minutes) 중심 수집 - 결과:
event.json
- 특정 이벤트 시점(
- CloudTrail
LookupEvents - AWS Config recorder 상태 확인
- AWS Config history 기반 리소스별 before/after 보강(best-effort)
현재 EVENT 계열은 아래 소스를 다룰 수 있도록 확장되어 있습니다.
MANUALSECURITYHUBAWS_HEALTH
특히 AWS Health 이벤트는 다음 정보를 extensions에 정리합니다.
extensions.event_originextensions.aws_healthextensions.actionability
이를 통해 “이 이벤트가 Terraform 조치 후보인지” 까지 Worker 단계에서 판단할 수 있습니다.
WEEKLY 실행 시 아래와 같은 운영 점검 결과를 extensions.advisor_checks[] 로 생성할 수 있습니다.
- 미사용 Elastic IP
- 미연결 EBS 볼륨
- RDS 백업 미설정
- RDS Multi-AZ 미설정
그리고 실행 상태와 요약은 아래에 기록됩니다.
extensions.advisor_collection_statusextensions.advisor_rollup
Worker는 로컬 산출물을 만든 뒤 아래처럼 S3에 업로드합니다.
raw/normalized/
예시:
s3://<bucket>/<prefix>/raw/...s3://<bucket>/<prefix>/normalized/canonical.jsons3://<bucket>/<prefix>/normalized/event.json
추가로 Worker는 raw/index.json 을 생성해서
이번 실행에서 남긴 evidence/normalized 파일 목록과 S3 URI를 inventory 형태로 기록합니다.
apps/worker/
dndn_worker/
__init__.py
run_job.py # Worker 핵심 실행 엔진
consumer.py # SQS payload consumer
s3_uploader.py # S3 업로드 유틸
Dockerfile # Worker consumer container image
tools/
run_payload.py # 로컬 payload 실행 도구
smoke_assume_role.py
smoke_cloudtrail.py
render_iam_templates.py
iam_templates/
customer_trust_policy.json
customer_permissions_policy.json
pyproject.toml
requirements.txt
README.md
Worker 핵심 로직입니다. 실제 수집/정규화/스키마 검증/S3 저장을 담당합니다. 공용 실행 진입점은 아래 두 함수입니다.
run_job_from_payload(payload: dict, ...)run_job_from_payload_file(path, ...)
로컬 job_dir 아래 산출물을 S3로 업로드하는 유틸입니다.
SQS 메시지 body를 payload JSON으로 받아 schema 검증 후
run_job_from_payload(...) 를 호출하는 운영용 consumer 진입점입니다.
python -m dndn_worker.consumer 를 기본 진입점으로 실행하는
Worker 컨테이너 이미지 정의입니다.
로컬에서 payload 파일 하나로 Worker를 실행할 때 사용합니다.
파일을 읽은 뒤 run_job_from_payload_file(...) 만 호출하는 개발/디버깅용 CLI 진입점입니다.
AssumeRole이 실제로 되는지 빠르게 확인하는 도구입니다.
CloudTrail 조회가 실제로 가능한지 빠르게 확인하는 도구입니다.
고객사 온보딩용 IAM 템플릿의 placeholder를 실제 값(principal ARN, external_id)으로 치환합니다.
고객사 계정 연동(AssumeRole)에 필요한 최소 IAM 템플릿입니다.
Worker는 항상 payload(JSON)를 입력으로 받습니다.
payload 스펙은 contracts/payload/job_payload.schema.json 을 따릅니다.
핵심 필드:
type:WEEKLY | EVENTaccount_idregionsassume_roles3- WEEKLY:
time_range - EVENT:
event_time,window_minutes,trigger
EVENT의 trigger 는 단순 메타데이터일 수도 있고,
상위 계층이 raw trigger event 자체를 넣어줄 수도 있습니다.
Worker는 가능한 경우 이 trigger 내용을 raw/trigger/*.json 으로 저장하고
meta.trigger.raw_event_s3_uri 를 결과에 남깁니다.
Worker는 실행 결과와 예외에 retryable 기준을 명시합니다.
- 실행 결과:
WorkerExecutionResult - 실행 예외:
WorkerExecutionError
원칙:
retryable=False이면 consumer는 메시지를 deleteretryable=True이면 consumer는 메시지를 남겨 재시도- validation 실패는
INVALID_PAYLOAD,retryable=False - S3 업로드 실패는
S3_PUT_FAILED,retryable=True
개발 단계에서는 assume_role.role_arn = "SELF" 로 두면,
AssumeRole 없이 현재 로컬 AWS 자격증명을 그대로 사용합니다.
이 모드는 다음 상황에서 유용합니다.
- 로컬 개발
- 스키마/정규화 테스트
- 실제 고객 계정 없이 구조 확인
같은 run_id가 다시 들어오면 로컬 out/<run_id>/normalized/{canonical|event}.json 존재 여부를 먼저 확인합니다.
- 이미 결과 파일이 있으면 재수집하지 않고
already_processed=True로 즉시 반환 - 즉, 동일 worker 인스턴스/볼륨 기준에서는
run_id재처리를 막습니다 - 분산 환경 전역 멱등성(S3/DB 락 기반)은 후속 운영 설계 범위입니다
운영/실전에서는 보통 고객 계정의 Read-only Role을 AssumeRole 해서 수집합니다.
Worker는 구조적으로 아래 두 세션을 분리합니다.
- collector session: 고객 계정에서 CloudTrail/Config를 읽는 세션
- storage session: 우리(DnDn) 계정 S3에 쓰는 세션
즉, 고객 계정 읽기와 우리 S3 저장을 분리해서 운영할 수 있게 되어 있습니다.
cd <repo-root>
python3 -m venv .venv
source .venv/bin/activate
pip install -U pip
pip install -e apps/workerpython - <<'PY'
import dndn_worker
print(dndn_worker.__file__)
PYaws sts get-caller-identity운영/consumer 경로와 로컬 실행 경로를 같은 코드로 맞추기 위해,
실제 job 실행은 dndn_worker.run_job.run_job_from_payload(...) 를 기준으로 두고
파일 기반 실행만 run_job_from_payload_file(...) / tools/run_payload.py 로 감쌉니다.
python apps/worker/tools/run_payload.py \
--payload /tmp/payload.weekly.json \
--repo-root . \
--out /tmp/dndn-out \
--max-events 500python apps/worker/tools/run_payload.py \
--payload /tmp/payload.event.json \
--repo-root . \
--out /tmp/dndn-out \
--max-events 200로컬에서는 대략 이런 구조로 생성됩니다.
/tmp/dndn-out/<run_id>/
raw/
index.json
meta/
job_payload.json
cloudtrail/
lookup_events.jsonl
event_<event_id>.json
config/
<resource_key>/
history.json
before.json
after.json
advisor/
...
trigger/
eventbridge.json
securityhub_finding.json
aws_health_event.json
normalized/
canonical.json # WEEKLY
event.json # EVENT
설명:
raw/는 evidence 원본 보관소입니다.normalized/는 report가 바로 읽을 표준 결과물입니다.raw/index.json은 이번 실행에서 어떤 파일이 생성됐는지 한 번에 보여주는 inventory 입니다.
queue message body는 완성된 payload JSON이라고 가정합니다.
python -m dndn_worker.consumer \
--queue-url https://sqs.ap-northeast-2.amazonaws.com/123456789012/dndn-worker \
--repo-root . \
--out /tmp/dndn-out \
--onceconsumer 동작:
- 메시지 수신
- JSON 파싱
- payload schema 검증
run_job_from_payload(...)호출visibility_timeout이 설정된 경우 처리 중 heartbeat로 메시지 visibility 연장retryable=False결과/예외면 deleteretryable=True예외면 delete 하지 않고 재시도 대상으로 남김- 같은
run_id재수신으로already_processed=True가 오면 delete
빌드:
docker build -f apps/worker/Dockerfile -t dndn-worker:local .실행 예시:
docker run --rm \
-e AWS_REGION=ap-northeast-2 \
-e DNDN_WORKER_QUEUE_URL=https://sqs.ap-northeast-2.amazonaws.com/123456789012/dndn-worker \
-e DNDN_WORKER_MAX_EVENTS=500 \
-e DNDN_WORKER_WAIT_TIME_SECONDS=20 \
-e DNDN_WORKER_MAX_MESSAGES=1 \
-v "$HOME/.aws:/home/worker/.aws:ro" \
dndn-worker:local기본 entrypoint:
python -m dndn_worker.consumer --repo-root /app --out /tmp/dndn-out한 번만 poll 하고 종료하려면:
docker run --rm \
-e AWS_REGION=ap-northeast-2 \
-e DNDN_WORKER_QUEUE_URL=https://sqs.ap-northeast-2.amazonaws.com/123456789012/dndn-worker \
-v "$HOME/.aws:/home/worker/.aws:ro" \
dndn-worker:local \
--onceconsumer 실행 시 주로 아래 env를 사용합니다.
DNDN_WORKER_QUEUE_URL: SQS queue URL. CLI--queue-url보다 기본값으로 사용DNDN_WORKER_MAX_EVENTS: job당 최대 CloudTrail event 수DNDN_WORKER_WAIT_TIME_SECONDS: SQS long polling wait timeDNDN_WORKER_MAX_MESSAGES: poll당 최대 수신 메시지 수DNDN_WORKER_HEARTBEAT_INTERVAL_SECONDS: 처리 중 visibility timeout 연장 heartbeat 주기.0이면visibility_timeout기준으로 자동 계산AWS_REGION또는AWS_DEFAULT_REGION: boto3 기본 리전- AWS credential chain 관련 env 또는 IAM role: 컨테이너/런타임에서 boto3가 사용하는 기본 인증 정보
현재 consumer는 worker의 retryable 기준에 맞춰 메시지 삭제 여부를 결정합니다.
retryable=False- 예:
INVALID_PAYLOAD,ASSUME_ROLE_FAILED - consumer는 메시지를 delete
- 예:
retryable=True- 예:
S3_PUT_FAILED, 일시적 AWS/network 오류 - consumer는 메시지를 남겨 재시도
- 예:
already_processed=True- 같은
run_id가 이미 처리된 상태 - consumer는 메시지를 delete
- 같은
운영 권장:
- SQS redrive policy(DLQ)로 최대 수신 횟수 제한
- visibility timeout은 평균 job 실행 시간보다 길게 설정
- job 실행 시간이 길거나 편차가 크면 heartbeat를 켜서 visibility timeout을 주기적으로 연장
DNDN_WORKER_MAX_MESSAGES=1부터 시작해 안정화 후 조정- retryable 실패는 CloudWatch/SQS metric 기반 알림 연결
고객 계정 AssumeRole 전 단계에서 Worker 기능 자체를 검증하려면, 현재 로그인된 AWS 계정을 수집 대상로 사용하고 S3도 DnDn 쪽 테스트 버킷으로 두는 방식이 가장 단순합니다.
검증 순서:
- 현재 AWS 자격증명으로 STS 호출이 되는지 확인
- 현재 AWS 자격증명으로 CloudTrail 조회가 되는지 확인
role_arn=SELFpayload로 WEEKLY 실행role_arn=SELFpayload로 EVENT 실행- 로컬 산출물과 S3 업로드 결과 확인
예시 명령:
cd /Users/mh/Desktop/DnDn-App
source .venv/bin/activate
export PYTHONPATH=apps/worker
python apps/worker/tools/smoke_assume_role.py --role-arn SELF
python apps/worker/tools/smoke_cloudtrail.py --role-arn SELF --region ap-northeast-2 --hours 24 --max 5
ACCOUNT_ID="$(aws sts get-caller-identity --query Account --output text)"
cat >/tmp/worker-weekly.json <<EOF
{
"account_id": "${ACCOUNT_ID}",
"assume_role": {
"external_id": "local-test",
"role_arn": "SELF"
},
"regions": ["ap-northeast-2"],
"rule_set_version": "eks-mvp-0.1",
"run_id": "manual-weekly-test-001",
"s3": {
"bucket": "dndn-data-dev-20260304",
"prefix": "manual-tests/weekly/manual-weekly-test-001"
},
"time_range": {
"start": "2026-03-01T00:00:00+09:00",
"end": "2026-03-08T00:00:00+09:00",
"timezone": "Asia/Seoul"
},
"type": "WEEKLY"
}
EOF
python apps/worker/tools/run_payload.py \
--payload /tmp/worker-weekly.json \
--repo-root . \
--out /tmp/dndn-out \
--max-events 20
cat >/tmp/worker-event.json <<EOF
{
"account_id": "${ACCOUNT_ID}",
"assume_role": {
"external_id": "local-test",
"role_arn": "SELF"
},
"event_time": "2026-03-11T10:30:00+09:00",
"hint": {
"resource": {
"region": "ap-northeast-2",
"resource_id": "my-eks-cluster",
"resource_type": "AWS::EKS::Cluster"
}
},
"regions": ["ap-northeast-2"],
"rule_set_version": "eks-mvp-0.1",
"run_id": "manual-event-test-001",
"s3": {
"bucket": "dndn-data-dev-20260304",
"prefix": "manual-tests/event/manual-event-test-001"
},
"trigger": {
"event_id": "manual-event-001",
"source": "EVENTBRIDGE"
},
"type": "EVENT",
"window_minutes": 60
}
EOF
python apps/worker/tools/run_payload.py \
--payload /tmp/worker-event.json \
--repo-root . \
--out /tmp/dndn-out \
--max-events 20
aws s3 ls s3://dndn-data-dev-20260304/manual-tests/weekly/manual-weekly-test-001/ --recursive
aws s3 ls s3://dndn-data-dev-20260304/manual-tests/event/manual-event-test-001/ --recursive이 테스트로 확인되는 것:
- 현재 AWS 자격증명 기반
SELF실행 가능 여부 - CloudTrail 실제 조회 가능 여부
- Worker 정규화 결과 생성 여부
- 로컬 raw/normalized 산출물 생성 여부
- DnDn S3 버킷 업로드 여부
이 테스트로 확인되지 않는 것:
- 고객 계정 AssumeRole trust policy
- 고객 계정 권한 정책
- 고객 환경별 리전/서비스 차이
즉, 이 절차는 "내 계정으로 Worker 자체가 실제로 도는지" 를 보는 실동작 테스트입니다.
- WEEKLY →
canonical.json - EVENT →
event.json
metacollection_statuseventsresourcesextensions
특히 meta.evidence 에는 다음 pointer가 들어갑니다.
raw_prefix_s3_urinormalized_prefix_s3_urijob_payload_s3_uriindex_s3_uri
각 단계의 수집 상태를 나타냅니다.
예:
assume_rolecloudtrailconfignormalized
상태 예:
OKNAFAILED
보고서 생성(B 파트)이 가장 많이 활용하는 단위입니다. 이벤트를 리소스 기준으로 묶어두고, 가능하면 Config 정보를 붙입니다.
contracts core를 깨지 않기 위해 확장 기능은 extensions 아래에 둡니다.
예:
event_originaws_healthactionabilityadvisor_collection_statusadvisor_checksadvisor_rollup
Worker evidence는 한 필드에 전부 모이지 않고 목적별로 나뉩니다.
- 실행 전체 artifact 위치:
meta.evidence - trigger 원문 위치:
meta.trigger.raw_event_s3_uri - CloudTrail raw 위치:
events[].raw - Config snapshot raw 위치:
resources[].config.before_s3_uri,after_s3_uri,extensions.history_s3_uri - advisor raw 위치:
extensions.advisor_checks[].evidence - Access Analyzer raw 위치:
extensions.access_analyzer_findings[].evidence.raw_s3_uri - Cost Explorer raw 위치:
extensions.cost_explorer_collection_status.ce.get_cost_and_usage.raw_s3_uri - CloudWatch raw 위치:
extensions.cloudwatch_alarms[].evidence.raw_s3_uri
사람이 한 번에 보기 가장 쉬운 entrypoint는 raw/index.json 입니다.
AWS Health 이벤트는 EventBridge payload를 기반으로 정규화됩니다.
Worker는 Health payload에서 다음을 읽습니다.
trigger.detail_typetrigger.logical_sourcetrigger.healthtrigger.resourcestrigger.health.affectedEntities
그리고 결과 JSON에 다음을 넣습니다.
extensions.event_origin.kind = AWS_HEALTHextensions.aws_healthextensions.actionability
Security Hub도 EVENT 루트에서 다룰 수 있도록 설계되어 있습니다. Worker는 trigger/finding 정보를 기반으로 resource ref를 보강하고, 필요한 경우 보고서/계획서 연결을 쉽게 할 수 있게 확장 필드를 유지합니다.
WEEKLY 실행에서는 Access Analyzer finding을 추가로 수집할 수 있습니다.
결과 JSON에는 아래 확장이 들어갑니다.
extensions.access_analyzer_collection_statusextensions.access_analyzer_findingsextensions.access_analyzer_rollup
이 finding은 외부 공개 가능 리소스, 교차 계정 접근 등 접근/권한 리스크를 보고서 단계에서 바로 활용할 수 있도록 요약됩니다.
WEEKLY 실행에서는 비용과 운영 경보 요약도 확장 필드로 수집할 수 있습니다.
결과 JSON에는 아래 확장이 들어갑니다.
extensions.cost_explorer_collection_statusextensions.cost_explorer_groupsextensions.cost_explorer_summaryextensions.cloudwatch_collection_statusextensions.cloudwatch_alarmsextensions.cloudwatch_rollup
이 확장은 주간 보고서에서 비용 변화와 현재 ALARM 상태를 함께 보여줄 때 사용합니다.
WEEKLY 실행 시 Worker는 주간 점검 항목을 수집해서 extensions.advisor_checks[]에 넣을 수 있습니다.
현재 기본 항목:
- 미사용 EIP
- 미연결 EBS
- RDS 백업 미설정
- RDS Multi-AZ 미설정
이 체크들은 “항상 결과가 있어야” 하는 건 아닙니다. 예를 들어 실제로 문제 리소스가 없으면:
advisor_checks = []advisor_rollup.total_checks = 0
이어도 정상입니다.
중요한 것은:
- 체크가 실행되었는지
advisor_collection_status가 채워졌는지- raw evidence가 남았는지 입니다.
Worker는 고객사 계정 연동을 위해 AssumeRole을 사용합니다.
이를 위해 iam_templates/ 와 render_iam_templates.py 를 제공합니다.
운영용 온보딩 절차와 권한 표는 별도 문서에 정리했습니다.
python apps/worker/tools/render_iam_templates.py \
--dndn-principal-arn arn:aws:iam::123456789012:role/DnDnWorkerRole \
--external-id dndn-tenant-abc출력:
out/iam_rendered/customer_trust_policy.rendered.jsonout/iam_rendered/customer_permissions_policy.rendered.json
- 고객 계정 Role은 읽기 전용 최소 권한
- 우리 S3 업로드는 고객 Role 권한이 아니라 DnDn 쪽 세션으로 수행
- ExternalId 기반 AssumeRole로 confused deputy 리스크를 줄임
대부분 editable install이 안 되어 있을 때 발생합니다.
해결:
pip install -e apps/workerCloudTrail/EventTime 등을 raw JSON으로 쓸 때 자주 발생합니다.
현재 Worker에는 _json_default()가 들어 있어 이 문제를 처리합니다.
계정/리전에 AWS Config recorder가 꺼져 있으면 정상입니다. 이 경우 Worker는 실패 대신:
NA(SERVICE_DISABLED)로 처리합니다.
고객 Role과 우리 S3 저장 세션이 분리되지 않았을 때 자주 생깁니다. 현재 Worker는 collector/storage session을 분리하는 구조를 사용합니다.
정상입니다. Worker는 raw evidence를 목적별 폴더에 분산 저장하고, normalized JSON에는 pointer만 남깁니다.
빠르게 확인하는 순서:
meta.evidence.index_s3_urimeta.trigger.raw_event_s3_urievents[].rawresources[].config
B는 Worker 결과 중 주로 아래를 사용합니다.
resources[]events[]extensions.aws_healthextensions.actionabilityextensions.advisor_checksresources[].config
C는 Worker 실행에 필요한 payload를 생성합니다. 핵심 필드:
typelinked account 정보(account_id / role_arn / external_id)time_range또는event_time / trigger
D는 EventBridge, IAM, 배포, 저장 구조와 연동합니다. 특히:
- AssumeRole principal ARN
- ExternalId
- S3 저장 구조 를 같이 맞춰야 합니다.
pip install -e apps/workeraws sts get-caller-identityrun_payload.py로 EVENT 또는 WEEKLY 실행
smoke_assume_role.pysmoke_cloudtrail.pyrun_payload.pywith role_arn
smoke_assume_role.py --role-arn SELFsmoke_cloudtrail.py --role-arn SELF- 위
6-8. 실제 AWS 계정으로 SELF 테스트절차 실행
docker build -f apps/worker/Dockerfile -t dndn-worker:local .DNDN_WORKER_QUEUE_URL설정python -m dndn_worker.consumer또는 Docker 실행 예시 사용
- WEEKLY payload 생성
run_payload.pyextensions.advisor_checks확인
이 Worker는 현재 DnDn에서 다음을 담당합니다.
- 변경 이력 수집
- 이벤트 보고서용 정규화
- 주간 보고서용 정규화
- AWS Health / SecurityHub 같은 이벤트 소스 보강
- 운영 점검(advisor checks)
- S3 저장
- AssumeRole 기반 고객 계정 수집