Skip to content

architecture iot data contract

minsoo edited this page Jun 9, 2026 · 3 revisions

IoT 데이터 계약

Edge data-plane이 AWS IoT Core로 전송하는 canonical JSON 계약이다. 코드 기준 source of truth는 apps/data-processor/processor/envelope.py, normalizer.py, apps/edge-iot-publisher/edge_iot_publisher.py다.


목적

이 계약은 edge 생성 컴포넌트와 cloud-side data processor 사이의 경계를 고정한다.

  • Edge는 원본 사실, 집계값, 업로드된 snapshot 참조만 보낸다.
  • 최종 위험 판단과 pipeline status 계산은 Lambda data-processor가 수행한다.
  • MQTT topic과 S3 raw partition은 factory_idsource_type만으로 분리한다.
  • adapter/generator/uploader는 outbox에 JSON 파일을 쓰고, publisher만 IoT Core에 publish한다.

Source Type

코드 기준 허용 source_type은 세 가지다.

factory_state
infra_state
image_snapshot
source_type 기본 생성 주기 의미
factory_state 3초 Risk Score 입력용 sensor/AI 집계
infra_state 20초 node/workload/device/heartbeat 상태
image_snapshot snapshot 파일 감지 시 S3에 업로드된 이미지의 metadata event

image_snapshot은 이미지 바이너리가 아니다. 이미지 파일은 presigned PUT으로 S3 image_snapshot/factory_id=... prefix에 먼저 올라가고, IoT Core에는 bucket/key/sha256/size 같은 metadata JSON만 전달된다.


MQTT Topic

aegis/{factory_id}/{source_type}

현재 사용 topic:

aegis/factory-a/factory_state
aegis/factory-a/infra_state
aegis/factory-a/image_snapshot
aegis/factory-b/factory_state
aegis/factory-b/infra_state
aegis/factory-c/factory_state
aegis/factory-c/infra_state

S3 Partition

IoT Rule raw key:

raw/{factory_id}/{source_type}/yyyy={YYYY}/mm={MM}/dd={DD}/{message_id}.json

Data processor processed key:

processed/{factory_id}/{dataset}/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json

datasetfactory_state, infra_state, image_snapshot, risk_score, state_snapshot 등이 될 수 있다.

이미지 바이너리 저장 key는 data processor raw/processed JSON과 분리된다.

image_snapshot/factory_id={factory_id}/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{filename}

예시:

raw/factory-a/image_snapshot/yyyy=2026/mm=06/dd=08/factory-a:image_snapshot:worker2:2026-06-08T09:42:35Z.json
processed/factory-a/image_snapshot/yyyy=2026/mm=06/dd=08/hh=09/factory-a:image_snapshot:worker2:2026-06-08T09:42:35Z.json
image_snapshot/factory_id=factory-a/yyyy=2026/mm=06/dd=08/hh=09/260608094235_event_FALLEN.jpg

Common Envelope

publish되는 메시지의 공통 envelope:

{
  "schema_version": "0.1.0",
  "message_id": "factory-a:factory_state:worker2:2026-05-14T01:00:00Z",
  "factory_id": "factory-a",
  "node_id": "worker2",
  "environment_type": "physical-rpi",
  "input_module_type": "sensor",
  "source_type": "factory_state",
  "source_timestamp": "2026-05-14T01:00:00Z",
  "published_at": "2026-05-14T01:00:01Z",
  "data_plane_instance_id": "edge-iot-publisher-7f8c9d",
  "payload": {}
}
필드 기준
schema_version 현재 0.1.0
message_id {factory_id}:{source_type}:{node_id}:{source_timestamp}
factory_id factory-a, factory-b, factory-c
node_id node 기준 메시지는 node ID, cluster 요약은 cluster
environment_type physical-rpi, vm-mac, vm-windows
input_module_type sensor, dummy, camera
source_type factory_state, infra_state, image_snapshot
source_timestamp 원본 데이터 기준 UTC ISO 8601
published_at publisher가 publish 직전에 UTC ISO 8601로 덮어씀
data_plane_instance_id publish한 Pod 또는 생성 컴포넌트 인스턴스 ID
payload source type별 JSON object

Lambda parser는 schema_version, message_id, factory_id, node_id, source_type, source_timestamp, published_at, payload 존재와 payload object 여부를 검증한다. edge-iot-publisher는 publish 전에 environment_type, input_module_type, data_plane_instance_id까지 포함한 확장 envelope를 검증한다.

message_id는 outbox 파일명, MQTT payload, S3 raw object key, processed body의 source_message_id, DynamoDB reference를 잇는 추적 키다.


factory_state

factory_state는 Risk Score 계산용 데이터다.

{
  "source_type": "factory_state",
  "payload": {
    "aggregation_window_seconds": 3,
    "sensor": {
      "sample_count": 5,
      "temperature_celsius_avg": 24.6,
      "humidity_percent_avg": 58.1,
      "pressure_hpa_avg": 1012.7
    },
    "ai_result": {
      "sample_count": 3,
      "fire_score": 0.0,
      "fall_score": 0.6667,
      "bend_score": 0.3333,
      "abnormal_sound": "none"
    }
  }
}

Normalizer 출력 필드:

출력 필드 입력
aggregation_window_seconds payload.aggregation_window_seconds
temperature_celsius payload.sensor.temperature_celsius_avg
humidity_percent payload.sensor.humidity_percent_avg
pressure_hpa payload.sensor.pressure_hpa_avg
sample_count payload.sensor.sample_count
fire_score payload.ai_result.fire_score
fall_score payload.ai_result.fall_score
bend_score payload.ai_result.bend_score
abnormal_sound payload.ai_result.abnormal_sound, 기본 none
ai_sample_count payload.ai_result.sample_count

Factory A는 InfluxDB 최근 window를 집계한다. Factory B/C는 profile 기반 synthetic 값을 만든다.


infra_state

infra_state는 운영 상태와 pipeline freshness 계산의 기준 데이터다.

{
  "source_type": "infra_state",
  "payload": {
    "heartbeat": {
      "agent_status": "alive",
      "last_successful_publish_at": null,
      "last_spool_write_status": "unknown",
      "last_spool_write_at": null,
      "publish_sequence": 123
    },
    "cluster": {
      "cluster_name": "factory-a",
      "kubernetes_version": "v1.34.6+k3s1"
    },
    "node_summary": { "total": 3, "ready": 3, "not_ready": 0 },
    "nodes": [
      {
        "node_id": "worker2",
        "role": "sensor-ai-audio-preferred",
        "ready": true,
        "cpu_usage_percent": 44.8,
        "memory_usage_percent": 63.0,
        "disk_usage_percent": 45.5,
        "network_reachability": "ok"
      }
    ],
    "workload_summary": { "total": 1, "running": 1, "not_running": 0 },
    "workloads": [
      {
        "namespace": "ai-apps",
        "name": "safe-edge-integrated-ai",
        "status": "Running",
        "ready": true,
        "containers_ready": 1,
        "containers_total": 1,
        "restart_count": 0,
        "node_id": "worker2"
      }
    ],
    "devices": {
      "bme280": { "available": true, "last_seen_at": "2026-05-14T01:00:00Z" },
      "camera": { "available": true, "last_seen_at": "2026-05-14T01:00:00Z" },
      "microphone": { "available": true, "last_seen_at": "2026-05-14T01:00:00Z" }
    }
  }
}

pipeline_heartbeat는 별도 source type이 아니다. Data processor는 최신 infra_state 시각과 현재 시각 차이로 pipeline_status를 계산하고, 1분 refresh schedule에서도 stale 상태를 재계산한다.


image_snapshot

image_snapshot은 업로드된 이미지 파일의 참조 이벤트다.

{
  "source_type": "image_snapshot",
  "payload": {
    "event_type": "FALLEN",
    "content_type": "image/jpeg",
    "size_bytes": 60345,
    "sha256": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
    "s3_bucket": "aegis-bucket-data",
    "s3_key": "image_snapshot/factory_id=factory-a/yyyy=2026/mm=06/dd=08/hh=09/260608094235_event_FALLEN.jpg",
    "local_path": "/var/lib/safe-edge/snapshots/260608094235_event_FALLEN.jpg",
    "upload_status": "uploaded"
  }
}

처리 순서:

  1. snapshot-uploader.jpg, .jpeg, .png snapshot 파일을 찾는다.
  2. presign Lambda에서 받은 URL로 이미지 바이너리를 S3에 PUT 한다.
  3. 업로드 성공 후 image_snapshot metadata JSON을 outbox에 쓴다.
  4. edge-iot-publisher가 metadata event를 aegis/factory-a/image_snapshot으로 publish한다.
  5. Data processor가 latest_image_snapshot, last_image_snapshot_at, processed JSON을 갱신한다.

Null / Missing Policy

영역 기준
envelope 필드 누락 parser 또는 publisher validation 실패. DynamoDB/processed에 반영하지 않음
payload 타입 반드시 JSON object
outbox 단계 published_at key는 존재해야 하며, publisher가 실제 publish 시각으로 덮어씀
센서 값 없음 sample_count: 0, normalizer 결과 수치 기본값은 0.0
AI 값 없음 sample_count: 0, score 기본값은 0.0, sound 기본값은 none
node/workload 수치 없음 null 유지
node/workload 배열 없음 빈 배열 [] 또는 summary 기반 계산
device available 없음 available: null, status: unknown
상태 문자열 없음 unknown 사용

필드 삭제보다 명시적 null, false, unknown, 빈 배열을 선호한다. 다만 현재 normalizer는 factory_state의 누락 수치와 score를 0.0으로 보정한다.


Local Spool/Outbox

adapter, dummy generator, snapshot uploader는 IoT Core에 직접 publish하지 않는다.

/var/lib/aegis/outbox/
/var/lib/aegis/outbox/tmp/
/var/lib/aegis/outbox/quarantine/
상황 처리
생성 tmp/에 쓰고 fsync 후 {message_id}.json으로 atomic rename
publish 성공 원본 .json 삭제
네트워크/MQTT 실패 파일 보존 후 재시도
JSON parse 실패 quarantine/ 이동
envelope validation 실패 quarantine/ 이동

edge-iot-publisher는 outbox root의 .json 파일만 대상으로 삼고, .snapshot-uploader-state.json과 하위 디렉터리는 무시한다.


관련 문서

Aegis-Pi Wiki

· 대표 문서 목록은 홈의 문서 탐색 표 참조

시작하기

요구사항

핵심 개념

아키텍처

컴포넌트 (Edge → Cloud → Dashboard)

Dashboard & 운영

시나리오 · 사례 · 참조

Clone this wiki locally