Metadata-Version: 2.4
Name: iris-vision-worker
Version: 1.2.1
Summary: Nedo Vision Worker Service Library for AI Vision Processing
Author-email: Willy Achmat Fauzi <willy.achmat@gmail.com>
Maintainer-email: Willy Achmat Fauzi <willy.achmat@gmail.com>
License-Expression: MIT
Project-URL: Homepage, https://gitlab.com/sindika/research/nedo-vision/nedo-vision-worker-service
Project-URL: Documentation, https://gitlab.com/sindika/research/nedo-vision/nedo-vision-worker-service/-/blob/main/README.md
Project-URL: Repository, https://gitlab.com/sindika/research/nedo-vision/nedo-vision-worker-service
Project-URL: Bug Reports, https://gitlab.com/sindika/research/nedo-vision/nedo-vision-worker-service/-/issues
Keywords: computer-vision,machine-learning,ai,worker-service,deep-learning,object-detection,neural-networks,video-processing
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: Operating System :: OS Independent
Classifier: Operating System :: POSIX :: Linux
Classifier: Operating System :: Microsoft :: Windows
Classifier: Operating System :: MacOS
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.8
Classifier: Programming Language :: Python :: 3.9
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Topic :: Multimedia :: Video
Classifier: Topic :: Scientific/Engineering :: Artificial Intelligence
Classifier: Topic :: System :: Hardware
Classifier: Environment :: GPU
Classifier: Environment :: No Input/Output (Daemon)
Requires-Python: >=3.8
Description-Content-Type: text/markdown
Requires-Dist: alembic>=1.8.0
Requires-Dist: ffmpeg-python>=0.2.0
Requires-Dist: grpcio>=1.50.0
Requires-Dist: pika>=1.3.0
Requires-Dist: protobuf>=6.31.1
Requires-Dist: psutil>=5.9.0
Requires-Dist: requests>=2.28.0
Requires-Dist: SQLAlchemy>=1.4.0
Requires-Dist: opencv-python>=4.6.0; platform_machine not in "aarch64 armv7l"
Requires-Dist: opencv-python-headless>=4.6.0; platform_machine in "aarch64 armv7l"
Requires-Dist: pynvml>=11.4.1; platform_system != "Darwin" or platform_machine != "arm64"
Provides-Extra: dev
Requires-Dist: pytest>=7.0.0; extra == "dev"
Requires-Dist: black>=22.0.0; extra == "dev"
Requires-Dist: isort>=5.10.0; extra == "dev"
Requires-Dist: mypy>=0.950; extra == "dev"
Requires-Dist: flake8>=4.0.0; extra == "dev"
Requires-Dist: pre-commit>=2.17.0; extra == "dev"

<div align="center">

# ai-vision-worker-service

**Agen edge Python IRIS — berjalan di hardware pabrik dekat kamera, menjalankan AI vision dan menstream deteksi ke manager**

Bagian dari ekosistem **IRIS — AI Vision Platform** PT Petrokimia Gresik

![Runtime](https://img.shields.io/badge/runtime-Python%203.8%2B%20(Docker%203.11)-3776AB)
![gRPC](https://img.shields.io/badge/gRPC-grpcio%201.50%2B%20%3A50051-244c5a)
![RabbitMQ](https://img.shields.io/badge/RabbitMQ-pika%201.3%2B%20%3A5672-FF6600)
![Cache](https://img.shields.io/badge/cache-SQLite%20%2B%20SQLAlchemy%2FAlembic-003B57)
![Version](https://img.shields.io/badge/version-1.2.1-blue)
![License](https://img.shields.io/badge/license-Proprietary-red)

</div>

---

## 1. Executive Summary

> **Untuk pembaca awam:** Bayangkan sebuah "otak kecil" yang dipasang di komputer industri persis di sebelah kamera CCTV di area pabrik. Otak kecil ini menyalakan program AI untuk menonton video kamera, mengenali pelanggaran keselamatan (mis. pekerja tanpa helm/APD, orang masuk area terlarang), lalu mengirim laporan ke server pusat. `ai-vision-worker-service` adalah otak kecil (agen) itu — bukan AI-nya sendiri, tapi *pengurus* yang mendaftar ke pusat, menerima perintah, menyalakan AI, dan mengirim hasil.

**Apa ini.** Sebuah **agen edge** (background daemon) berbasis Python yang dijalankan di mesin dekat kamera pabrik. Ia tidak melakukan inference sendiri — inference dikerjakan library **`ai-vision-worker-core`** — melainkan bertugas sebagai *control plane* di sisi edge: registrasi & autentikasi ke manager, menerima perintah start/stop pipeline, menyalakan/mematikan worker-core, mem-buffer hasil deteksi ke database lokal, lalu meng-upload-nya ke pusat secara batch. Juga melapor metrik sistem (CPU/RAM/GPU/latency) secara berkala.

**Masalah yang dipecahkan.** Kamera pabrik tersebar dan koneksi ke pusat tidak selalu stabil. Kalau semua video dikirim mentah ke pusat, bandwidth jebol dan latency tinggi. Solusinya: **proses di edge, kirim hanya kesimpulan.** Agen ini memungkinkan AI berjalan lokal (dekat kamera, hemat bandwidth, tahan putus jaringan karena hasil di-buffer di SQLite dulu) sambil tetap dikendalikan terpusat oleh `ai-vision-manager`.

**Peran dalam ekosistem IRIS.**

| Arah | Service tetangga | Hubungan |
|---|---|---|
| Upstream (dikendalikan oleh) | `ai-vision-manager` (.NET) | gRPC `:50051` (registrasi, upload deteksi, status) + RabbitMQ `:5672` (terima perintah) |
| Sibling (library AI) | `ai-vision-worker-core` (Python) | Berbagi `--storage-path` (SQLite + model) — worker-core menulis deteksi, agen ini meng-upload |
| Downstream (video preview) | MediaMTX / RTMP `:1935` | Publish preview annotated ke RTMP untuk ditonton di frontend |
| Broker | RabbitMQ (exchange `nedo.*`) | Direct exchange, routing key = `worker_id` |

### Status saat ini

| Aspek | Nilai |
|---|---|
| Versi | `1.2.1` (`nedo_vision_worker/__init__.py`) |
| Bahasa/runtime | Python 3.8+ (image produksi `python:3.11-slim`) |
| LOC (Python) | ± 15.088 baris di `nedo_vision_worker/` |
| Port keluar | gRPC `:50051` (ke manager), AMQP `:5672` (RabbitMQ), RTMP `:1935` (preview) |
| Dependensi kunci | grpcio ≥1.50, pika ≥1.3, protobuf ≥6.31.1, SQLAlchemy ≥1.4, alembic ≥1.8, psutil ≥5.9, opencv-python ≥4.6, pynvml ≥11.4.1 |
| Registry | GHCR `ghcr.io/tekinfopg/*`, auto-deploy via Watchtower |
| Maturity (CMMI-style) | **Level 2 — Managed** (menuju 3) |

### Stack ringkas

```
Runtime      │ Python 3.8+ · daemon proses tunggal · multithread (bukan async)
Transport    │ gRPC (grpcio) :50051 — kontrol & upload; RabbitMQ (pika) :5672 — perintah
Auth         │ Token worker (dari frontend) → metadata Bearer (TokenAuthInterceptor)
Cache lokal  │ SQLite (default.db / config.db / logging.db) via SQLAlchemy + Alembic
AI           │ ai-vision-worker-core (library terpisah, shared storage-path)
Sistem       │ psutil (CPU/RAM) · pynvml (GPU NVIDIA) · ffmpeg (video)
Video        │ FFmpeg → RTMP :1935 (preview) · RTSP :8554 (probe sumber)
Deploy       │ Docker (python:3.11-slim) · GHCR · Watchtower · self-hosted CI
```

---

## 2. Proses Bisnis (BPMN)

> **Untuk awam & analis:** Diagram di bawah menceritakan "hidup satu worker" dari komputer edge dinyalakan sampai deteksi mengalir ke pusat. Perhatikan tiga pelaku: **Perangkat Edge** (komputer di pabrik), **Agen** (program ini), dan **Manager** (server pusat).

```mermaid
flowchart TB
    subgraph EDGE["🖥️ Perangkat Edge (pabrik)"]
        Boot([Device boot / container start]):::start
        HWID[Ambil Hardware ID unik]
        CoreProc[worker-core jalan → inference kamera]
    end

    subgraph AGENT["🤖 Agen ai-vision-worker-service"]
        Auth[Autentikasi token ke manager]
        CheckCfg{Config lokal<br/>sudah ada?}
        Register[Registrasi: GetConnectionInfo<br/>→ terima worker_id + kredensial RabbitMQ]
        SaveCfg[Simpan config ke SQLite]
        Threads[Nyalakan worker threads<br/>sync · sender · listener]
        WaitCmd[Dengarkan perintah RabbitMQ]
        Gate{Perintah?}
        RunCore[Start processing workers<br/>→ status RUN ke manager]
        StopCore[Stop processing workers<br/>→ status STOP]
        Buffer[Baca deteksi worker-core<br/>dari SQLite bersama]
        Upload[Batch-upload deteksi + gambar via gRPC]
    end

    subgraph MGR["☁️ ai-vision-manager (pusat)"]
        Assign[Assign pipeline ke worker]
        Store[(PostgreSQL + storage)]
        Metrics[Terima heartbeat metrik sistem]
    end

    Boot --> HWID --> Auth
    Auth --> CheckCfg
    CheckCfg -- belum --> Register --> SaveCfg --> Threads
    CheckCfg -- sudah --> Threads
    Register -.token.-> Assign
    Threads --> WaitCmd --> Gate
    Assign -.AMQP command.-> Gate
    Gate -- start --> RunCore --> CoreProc
    Gate -- stop --> StopCore
    CoreProc --> Buffer --> Upload -.gRPC.-> Store
    Threads -.heartbeat 10-30s.-> Metrics
    RunCore -.UpdateStatus.-> Assign

    classDef start fill:#c8e6c9,stroke:#2e7d32
```

### Tabel langkah proses

| No | Aktivitas | Aktor | Sistem/Tool | Output |
|----|-----------|-------|-------------|--------|
| 1 | Boot perangkat / start container | Edge | systemd / Docker | Proses agen hidup |
| 2 | Ambil Hardware ID unik (UUID) | Agen | `HardwareID.get_unique_id()` | `device_id` |
| 3 | First-time setup: kirim token → minta info koneksi | Agen → Manager | gRPC `GetConnectionInfo` | `worker_id` + kredensial RabbitMQ |
| 4 | Simpan konfigurasi lokal | Agen | SQLite `config.db` | Config persisted |
| 5 | Nyalakan worker threads | Agen | `WorkerManager.start_all()` | 8 thread aktif |
| 6 | Dengarkan perintah pipeline | Agen | RabbitMQ direct exchange | Consumer siap |
| 7 | Manager assign & kirim `start` | Manager → Agen | AMQP `nedo.worker.core.action` | Processing dimulai |
| 8 | Jalankan inference | Agen → worker-core | shared `--storage-path` | Deteksi ditulis ke SQLite |
| 9 | Buffer & batch-upload deteksi | Agen → Manager | gRPC `UpsertBatch` / `SendPipelineDetectionData` | Deteksi tersimpan di pusat |
| 10 | Lapor metrik sistem berkala | Agen → Manager | gRPC `SendSystemUsage` | Heartbeat CPU/RAM/GPU |

---

## 3. Arsitektur

### 3.1 Komponen high-level

> **Awam:** kotak tengah adalah program ini. Panah menunjukkan "siapa bicara ke siapa, lewat jalur apa". gRPC = jalur cepat biner untuk kontrol & data; AMQP = jalur pesan/perintah; RTMP = jalur video.

```mermaid
flowchart LR
    subgraph Edge["🏭 Mesin Edge (dekat kamera)"]
        CAM["📹 Kamera<br/>RTSP"]
        subgraph SVC["🤖 ai-vision-worker-service"]
            WM["WorkerManager<br/>(8 thread)"]
            GCM["GrpcClientManager<br/>(singleton clients)"]
            SQL[("SQLite<br/>default/config/logging")]
        end
        CORE["🧠 ai-vision-worker-core<br/>(YOLO · RF-DETR · SFSORT)"]
        MTX["📡 MediaMTX / RTMP"]
    end

    subgraph Central["☁️ Pusat IRIS"]
        MGR["ai-vision-manager<br/>(.NET)"]
        MQ["RabbitMQ<br/>exchange nedo.*"]
        PG[("PostgreSQL")]
    end

    CAM -->|RTSP :8554| CORE
    CORE -->|"tulis deteksi + gambar<br/>(shared storage-path)"| SQL
    SVC -->|"baca buffer"| SQL
    CORE -->|annotated| MTX
    MTX -->|RTMP :1935| Central

    WM <-->|"AMQP :5672<br/>perintah (by worker_id)"| MQ
    MQ --- MGR
    GCM -->|"gRPC :50051<br/>register · upload · status · metrik"| MGR
    MGR --> PG

    style SVC fill:#244c5a,color:#fff
    style CORE fill:#6a1b9a,color:#fff
    style MGR fill:#1565c0,color:#fff
```

### 3.2 Sequence — boot → auth → registrasi → deteksi

> **Awam:** urutan waktu satu worker dari nyala sampai deteksi pertama sampai ke pusat.

```mermaid
sequenceDiagram
    autonumber
    participant OS as 🖥️ Edge OS
    participant SVC as 🤖 WorkerService
    participant MGR as ☁️ Manager (gRPC :50051)
    participant MQ as 🐇 RabbitMQ (:5672)
    participant CORE as 🧠 worker-core
    participant DB as 🗄️ SQLite

    OS->>SVC: python -m nedo_vision_worker.cli run --token XXX
    SVC->>SVC: set_storage_path() + init SQLite
    SVC->>SVC: HardwareID.get_unique_id() → device_id
    alt Config belum ada (first-time)
        SVC->>MGR: GetConnectionInfo(token)
        MGR-->>SVC: worker_id + rabbitmq host/port/user/pass
        SVC->>DB: simpan config (config.db)
    end
    SVC->>SVC: WorkerManager.start_all() (8 thread)
    SVC->>MQ: consume exchange nedo.worker.core.action (rk=worker_id)
    MGR->>MQ: publish {action:"start"} (rk=worker_id)
    MQ-->>SVC: perintah start
    SVC->>CORE: _start_workers() → video/pipeline/sender aktif
    SVC->>MGR: UpdateStatus("run")
    loop tiap frame
        CORE->>DB: tulis deteksi (default.db) + gambar
    end
    loop tiap ~10s
        SVC->>DB: baca batch deteksi
        SVC->>MGR: UpsertBatch / SendPipelineDetectionData
        SVC->>MGR: SendSystemUsage(CPU/RAM/GPU/latency, version)
    end
```

### 3.3 Katalog gRPC call ke manager

> **Spesialis:** semua RPC via `VisionWorkerServiceStub` dan stub deteksi khusus. Channel dibuat `make_grpc_channel()` — insecure by default, TLS bila `GRPC_TLS_ENABLED=true`. Token dikirim **dua jalur**: di metadata `authorization: Bearer <token>` (`TokenAuthInterceptor`) dan di body request (back-compat). `GrpcClientBase` punya circuit breaker (buka setelah 3 gagal, cooldown 15s), deadline per-call 10s, dan reconnect di background.

| Service (proto) | RPC | Dipakai oleh | Fungsi |
|---|---|---|---|
| `VisionWorkerService` | `GetConnectionInfo` | `ConnectionInfoClient` | Registrasi first-time: tukar token → worker_id + kredensial RabbitMQ |
| `VisionWorkerService` | `SendSystemUsage` | `SystemUsageClient` | Heartbeat metrik (CPU, RAM, GPU[], latency_ms, version) |
| `VisionWorkerService` | `UpdateStatus` | `WorkerStatusClient` | Lapor status worker (`run`/`stop`) |
| `HealthCheckService` | `HealthCheck` | `Networking.check_grpc_latency` | Ukur latency gRPC (tiap 10s) |
| `ImageService` | `GetLastImageDate` / `UploadImage` | `ImageUploadClient` | Upload frame/gambar bukti pelanggaran |
| `WorkerSourceService` | `GetWorkerSourceList` / `Update` / `DownloadSourceFile` | `WorkerSourceClient` | Sinkron daftar sumber kamera + unduh file sumber |
| `WorkerSourcePipelineService` | `GetListByWorkerId`, `SendPipelineImage`, `UpdateStatus`, `SendPipelineDebug`, `SendPipelineDetectionData` | `WorkerSourcePipelineClient` | Sinkron pipeline + upload deteksi pipeline + debug |
| `AIModelGRPCService` | `GetAIModelList` / `DownloadAIModel` | `AIModelClient` | Sinkron & unduh model AI ke storage bersama |
| `PPEDetectionGRPCService` | `Upsert` / `UpsertBatch` | `PPEDetectionClient` | Batch-upload deteksi APD/PPE |
| `ClumpDetectionGRPCService` | `UpsertBatch` | `ClumpDetectionClient` | Batch-upload deteksi gumpalan (clump) |
| `SpeedDetectionGRPCService` | `UpsertBatch` | `SpeedDetectionClient` | Batch-upload deteksi kecepatan |
| `ColorAnomalyDetectionGRPCService` | `UpsertBatch` | `ColorAnomalyDetectionClient` | Batch-upload anomali warna |
| `RestrictedArea` (via pipeline) | violation batch | `RestrictedAreaClient` | Batch-upload pelanggaran area terlarang |
| `DatasetSourceService` | frame upload | `DatasetSourceClient` | Kirim frame untuk pengumpulan dataset training |
| `VideoStreamService` | `StreamVideo` (stream) | `VideoStreamClient` | Streaming frame ke manager (preview alternatif) |

### 3.4 Katalog perintah RabbitMQ (inbound)

> **Awam:** manager tidak "menelepon" tiap worker langsung — ia menaruh pesan ke kotak surat (exchange) dengan alamat = `worker_id`. Hanya worker yang alamatnya cocok yang membaca. Semua exchange bertipe **direct** dan `durable`; queue worker `auto_delete` + `exclusive`; routing key = `worker_id` (huruf kecil).

| Exchange | Queue | Worker konsumen | Payload → Aksi |
|---|---|---|---|
| `nedo.worker.core.action` | `nedo.worker.core.{worker_id}` | `CoreActionWorker` | `{action}`: `start`/`stop`/`restart`/`debug` → nyalakan/matikan seluruh processing worker |
| `nedo.worker.pipeline.action` | `nedo.worker.pipeline.{worker_id}` | `PipelineActionWorker` | `{workerSourcePipelineId, action}` → set `pipeline_status_code` per pipeline (`run`/`stop`/`restart`), `debug` → buat debug entry |
| `nedo.pipeline.image.request` | `nedo.pipeline.request.{worker_id}` | `PipelineImageWorker` | Permintaan snapshot gambar pipeline |
| `nedo.worker.stream.preview` | `nedo.worker.preview.{worker_id}` | `VideoStreamWorker` | Trigger preview video (RTMP) |
| `nedo.worker.source.probe.request` | `nedo.worker.probe.{worker_id}` | `StreamProbeWorker` | `{requestId, url}` → probe sumber (RTSP/file) berlapis, balas ke `nedo.worker.source.probe.response` (rk=`requestId`) |

Koneksi RabbitMQ: `pika.SelectConnection` (event loop), heartbeat 30s, **reconnect eksponensial** (5s → max 60s), `prefetch_count=1`, `auto_ack=True`. Isolasi environment lewat `RABBITMQ_VHOST` (default `/`) — prod & staging bisa berbagi satu broker korporat tanpa saling silang.

### 3.5 Hubungan dengan `ai-vision-worker-core`

> **Awam:** agen ini dan library AI adalah dua proses terpisah yang "berbagi laci" — satu folder di disk. AI menulis hasil ke laci, agen mengambilnya dan mengirim ke pusat.

- **Kontrak berbagi:** keduanya harus dijalankan dengan **`--storage-path` yang sama**. Isi laci: `sqlite/` (default.db deteksi, config.db, logging.db), `files/detection_image/`, `model/`, dan `.worker_core_version`.
- **worker-core → menulis:** hasil inference (PPE, clump, speed, color anomaly, restricted-area) + gambar bukti ke `default.db` dan `files/detection_image/`.
- **worker-service → membaca & mengunggah:** `DataSenderWorker` menjalankan `*DetectionManager.send_*_batch()` tiap ~10s membaca buffer SQLite (batch di-cap agar backlog saat manager down tidak membanjiri RAM), upload via gRPC, lalu menghapus baris yang sudah tersinkron.
- **Versi gabungan:** `get_worker_version()` membaca `.worker_core_version` dari storage bersama → laporkan `"svc 1.2.1 / core X.Y.Z"` di field `version` heartbeat. Jadi manager/UI tahu versi keduanya walau prosesnya terpisah.

### 3.6 ERD — cache SQLite lokal

```mermaid
erDiagram
    CONFIG_DB ||--|| SERVER_CONFIG : "key-value"
    SERVER_CONFIG {
        string key PK "worker_id, server_host, token, rabbitmq_*"
        string value
    }
    DEFAULT_DB ||--o{ PPE_DETECTION : buffers
    DEFAULT_DB ||--o{ RESTRICTED_AREA_VIOLATION : buffers
    DEFAULT_DB ||--o{ CLUMP_DETECTION : buffers
    DEFAULT_DB ||--o{ SPEED_DETECTION : buffers
    DEFAULT_DB ||--o{ COLOR_ANOMALY_DETECTION : buffers
    DEFAULT_DB ||--o{ WORKER_SOURCE_PIPELINE_DETECTION : buffers
    WORKER_SOURCE_PIPELINE_DETECTION {
        string id PK
        string worker_source_pipeline_id FK
        string detection_image "path di files/detection_image"
        datetime created_at "drained batch-by-batch"
    }
```

Tiga database SQLite terpisah (`DatabaseManager`): `default.db` (buffer deteksi & sumber), `config.db` (konfigurasi), `logging.db` (log). Skema di-auto-migrate saat startup lewat Alembic autogenerate (`produce_migrations`), dengan penanganan khusus SQLite untuk drop-NOT-NULL.

### 3.7 Security

```mermaid
flowchart LR
    T["🔑 Token worker<br/>(dari frontend)"] --> I["TokenAuthInterceptor"]
    I -->|"metadata: authorization Bearer ***"| CH["gRPC channel"]
    CH -->|"GRPC_TLS_ENABLED=true?"| TLS{TLS}
    TLS -- ya --> SEC["secure_channel<br/>(CA: GRPC_TLS_CA_PATH)"]
    TLS -- tidak --> INS["insecure_channel"]
    SEC --> MGR["Manager gRPC auth interceptor"]
    INS --> MGR
    MGR -- UNAUTHENTICATED --> CB["_notify_auth_failure()<br/>→ service shutdown"]

    style T fill:#fff3e0
    style CB fill:#ffcdd2
```

- **Token redaction:** semua log melewati regex `_redact()` yang menyembunyikan `authorization/bearer/token` sebelum ditulis. Config `token`/`password` di-mask (`***`) saat di-print.
- **Auth-failure hard stop:** RPC `UNAUTHENTICATED`/`PERMISSION_DENIED` → `set_auth_failure_callback` men-trigger `WorkerService.stop()` + exit code 1 (worker mati kalau token dicabut).
- **TLS opsional** via env; default insecure (cocok untuk deploy host langsung di jaringan internal pabrik, gRPC 50051 direct host).

### 3.8 Ports & Protokol

| Port | Protokol | Arah | Tujuan | Tool/Library |
|---|---|---|---|---|
| `50051` | gRPC/HTTP2 | keluar → manager | Register, upload deteksi, status, metrik, health | grpcio |
| `5672` | AMQP | keluar → RabbitMQ | Terima perintah pipeline (direct, by worker_id) | pika |
| `1935` | RTMP | keluar → MediaMTX | Publish preview video annotated | ffmpeg |
| `8554` | RTSP | masuk ← kamera | Sumber video (di-probe & dikonsumsi worker-core) | ffmpeg/OpenCV |

> Port `50051`, `1935`, `8554` di-`EXPOSE` di Dockerfile. Untuk agen edge, koneksi bersifat *outbound* ke pusat — tidak ada port inbound yang perlu dibuka dari internet.

### 3.9 Tech Stack & Rationale

| Pilihan | Versi | Alasan |
|---|---|---|
| gRPC (grpcio) | ≥1.50 | RPC biner cepat & hemat bandwidth untuk kontrol + upload deteksi dari edge; streaming untuk video |
| protobuf | ≥6.31.1 | Stub `*_pb2` di-generate dengan gencode 6.31.x — runtime wajib match (lihat commit fix Speed/Clump) |
| pika | ≥1.3 | Klien RabbitMQ murni-Python; `SelectConnection` untuk consumer event-loop non-blocking |
| SQLAlchemy + Alembic | ≥1.4 / ≥1.8 | ORM + auto-migrate SQLite; buffer deteksi tahan-putus di edge |
| psutil / pynvml | ≥5.9 / ≥11.4.1 | Metrik CPU/RAM (semua platform) & GPU NVIDIA (Jetson/dGPU); pynvml di-skip di Apple Silicon |
| opencv-python(-headless) | ≥4.6 | Var headless otomatis untuk ARM/aarch64 (Jetson, RPi) |
| ffmpeg-python | ≥0.2 | Wrapper FFmpeg untuk RTSP→RTMP & probe sumber |

---

## 4. Tata Kelola & Kematangan

> **Untuk manajemen & auditor:** pemetaan repo ke kerangka tata kelola TI. Bukti diambil dari artefak nyata di repo (CI, kode, commit), bukan klaim.

### COBIT 2019

| Objective | Bagaimana repo memenuhinya |
|---|---|
| **APO03** Managed Enterprise Architecture | Peran edge-agent terdefinisi jelas dalam ekosistem IRIS; batas tanggung jawab (kontrol vs inference) dipisah dari `worker-core` |
| **BAI03** Managed Solutions Build | Struktur berlapis (services/worker/repositories/models); proto sebagai kontrak antar-service; `README_DEV.md` untuk regen protobuf |
| **BAI06** Managed IT Changes | GitHub Actions CI (`ci.yml`) di `develop`/`main`; commit terstruktur (`fix(grpc):`, `feat(rabbitmq):`); Watchtower untuk rollout terkontrol |
| **DSS01** Managed Operations | Auto-reconnect RabbitMQ (backoff eksponensial) + circuit breaker gRPC + graceful shutdown (SIGINT/SIGTERM); `doctor` untuk pre-flight check |
| **DSS05** Managed Security Services | Token redaction di log, auth-failure hard-stop, TLS opsional, kredensial di-mask; vhost isolation antar-env |
| **MEA01** Performance Monitoring | Heartbeat metrik sistem (CPU/RAM/GPU/latency) tiap 10-30s ke manager; version reporting gabungan svc+core |

### PMBOK / Knowledge Area

| Area | Deliverable konkret di repo |
|---|---|
| Scope | CLI `run`/`doctor` dengan argumen terdefinisi; peran agen (bukan AI) eksplisit |
| Schedule | Interval terkonfigurasi (`--system-usage-interval`, sync 10s, sender 10s) |
| Quality | CI: `ruff check` + `compileall` + `pytest`; `tests/` (proto parity, stream probe); mypy strict config di `pyproject.toml` |
| Risk | Circuit breaker, deadline per-call 10s, batch cap anti-OOM, backoff reconnect, buffer SQLite tahan-putus |
| Integration | Proto sebagai kontrak gRPC; exchange `nedo.*` sebagai kontrak AMQP; shared `storage-path` dengan worker-core |

### IT Maturity (CMMI-style)

**Level saat ini: 2 — Managed (menuju 3 Defined).**

Justifikasi: proses inti sudah *managed* — build terulang (Docker + CI), penanganan error/reconnect matang, keamanan token diperhatikan, observability via heartbeat, dan ada test otomatis. Namun untuk naik ke **Level 3 (Defined)** masih kurang: (1) coverage test tipis (hanya 2 file test, logika murni; tidak ada integration test gRPC/RabbitMQ), (2) `ruff check` masih `continue-on-error` (belum gate keras), (3) konfigurasi tersebar di SQLite runtime tanpa `.env.example` terdokumentasi, (4) belum ada runbook operasional/observability terpusat (metrik hanya ke manager, bukan Prometheus). Isu operasional seperti *detection-stall* (lihat §9) belum terinstrumentasi dengan alert otomatis.

---

## 5. Repository Structure

```
ai-vision-worker-service/
├── nedo_vision_worker/
│   ├── cli.py                     # Entrypoint CLI: subcommand run / doctor
│   ├── worker_service.py          # WorkerService — lifecycle, first-time setup, signal handling
│   ├── doctor.py                  # Diagnostik sistem (ffmpeg, opencv, gpu, storage)
│   ├── initializer/
│   │   └── AppInitializer.py      # Registrasi first-time: token → GetConnectionInfo → simpan config
│   ├── worker/                    # 🧵 Thread workers (dikelola WorkerManager)
│   │   ├── WorkerManager.py       #   Orkestrasi 8 worker; gate start/stop processing
│   │   ├── CoreActionWorker.py    #   Listener RabbitMQ core.action (start/stop/restart)
│   │   ├── PipelineActionWorker.py#   Listener pipeline.action (per-pipeline status)
│   │   ├── DataSyncWorker.py      #   Sync model/source/pipeline/detection dari manager (10s)
│   │   ├── DataSenderWorker.py    #   Batch-upload deteksi + metrik + gambar (10s)
│   │   ├── VideoStreamWorker.py   #   Preview video RTMP
│   │   ├── PipelineImageWorker.py #   Snapshot gambar pipeline on-demand
│   │   ├── DatasetFrameWorker.py  #   Kirim frame untuk dataset training
│   │   ├── StreamProbeWorker.py   #   Probe sumber RTSP/file (req/resp RabbitMQ)
│   │   ├── SystemUsageManager.py  #   Kumpul & kirim metrik + latency thread
│   │   ├── RabbitMQListener.py    #   Consumer pika reusable (reconnect backoff)
│   │   └── *DetectionManager.py   #   PPE/Clump/Speed/ColorAnomaly/RestrictedArea batch senders
│   ├── services/                  # 🔌 gRPC clients (singleton via GrpcClientManager)
│   │   ├── GrpcClientBase.py      #   Base: circuit breaker, retry, redaction, auth-failure
│   │   ├── GrpcClientManager.py   #   Singleton pool client (reuse channel)
│   │   ├── ConnectionInfoClient.py#   GetConnectionInfo (registrasi)
│   │   ├── WorkerStatusClient.py  #   UpdateStatus
│   │   ├── SystemUsageClient.py   #   SendSystemUsage
│   │   ├── AIModelClient.py       #   Sync & unduh model
│   │   └── ...                    #   Video/Image/Dataset/Detection clients + RTMP streamers
│   ├── repositories/              # 💾 Data access SQLite (SQLAlchemy)
│   ├── models/                    #   ORM entities (config, detection, pipeline, ...)
│   ├── protos/                    # 📐 .proto + stub *_pb2 / *_pb2_grpc (gencode 6.31.1)
│   ├── config/ConfigurationManager.py  # CRUD config di SQLite
│   ├── database/DatabaseManager.py     # Engine/session SQLite + auto-migrate
│   └── util/                      # HardwareID, Networking(TLS/interceptor), SystemMonitor, Version, PlatformDetector
├── tests/                         # test_proto_schema_parity, test_stream_probe
├── Dockerfile                     # python:3.11-slim + ffmpeg; CMD python -m ...cli
├── docker-compose.local.yml       # Dev: join nedo-network, shared /app/data dengan core
├── requirements.txt / pyproject.toml
├── install.sh / install.bat / run.sh / run.bat
├── PLATFORM_SUPPORT.md            # Matriks platform (Linux/Win/macOS/Jetson/ARM)
└── .github/workflows/             # ci.yml, docker-build-and-push.yml
```

---

## 6. Konfigurasi & Environment

Konfigurasi utama lewat **argumen CLI**; sebagian kecil lewat **environment variable**. Kredensial RabbitMQ **tidak** di-set manual — didapat otomatis dari manager saat registrasi dan disimpan di SQLite `config.db`.

### Argumen CLI (`run`)

| Argumen | Wajib | Default | Fungsi |
|---|---|---|---|
| `--token` | ✅ | — | Token autentikasi worker (dari frontend) |
| `--server-host` | ❌ | `be.vision.sindika.co.id` | Host gRPC manager |
| `--server-port` | ❌ | `50051` | Port gRPC manager |
| `--rtmp-server` | ❌ | `rtmp://live.vision.sindika.co.id:1935/live` | Target RTMP preview |
| `--storage-path` | ❌ | `data` | **Wajib sama dengan worker-core** — folder SQLite/model bersama |
| `--system-usage-interval` | ❌ | `30` | Interval lapor metrik (detik) |
| `--log-level` | ❌ | `INFO` | DEBUG/INFO/WARNING/ERROR/CRITICAL |

### Environment variables

| Nama | Wajib | Default | Fungsi |
|---|---|---|---|
| `RABBITMQ_VHOST` | ❌ | `/` | Isolasi broker per-environment (prod vs staging berbagi 1 broker) |
| `GRPC_TLS_ENABLED` | ❌ | `false` | Aktifkan `secure_channel` untuk gRPC |
| `GRPC_TLS_CA_PATH` | ❌ | — | Path CA root untuk verifikasi TLS |
| `PYTHONUNBUFFERED` | ❌ | — | Log real-time di container (di-set di compose) |

### Konfigurasi tersimpan (SQLite `config.db`)

Diisi otomatis saat registrasi: `worker_id`, `server_host`, `server_port`, `token`, `rabbitmq_host`, `rabbitmq_port`, `rabbitmq_username`, `rabbitmq_password`, `rabbitmq_vhost`.

---

## 7. Local Development

### Prasyarat

```bash
python --version    # 3.8+ (3.11 disarankan, sesuai image)
ffmpeg -version     # wajib (RTSP/RTMP)
docker --version    # 24+ (opsional, untuk container)
```

### Setup

```bash
# Clone
git clone https://github.com/tekinfopg/ai-vision-worker-service
cd ai-vision-worker-service

# Virtualenv + install
python -m venv venv && source venv/bin/activate   # Windows: venv\Scripts\activate
pip install -e .            # atau: pip install -r requirements.txt

# Cek kesiapan sistem
python -m nedo_vision_worker.cli doctor

# Jalankan (butuh manager + RabbitMQ + token)
python -m nedo_vision_worker.cli run \
  --token YOUR_TOKEN \
  --server-host localhost --server-port 50051 \
  --storage-path ./data \
  --rtmp-server rtmp://localhost:1935/live
```

### Docker (dev, `docker-compose.local.yml`)

```bash
# Prasyarat: network 'nedo-network' + manager/rabbitmq/mediamtx sudah jalan
WORKER_TOKEN=xxxxx docker compose -f docker-compose.local.yml up --build
```

Compose memakai `--server-host manager` dan `--rtmp-server rtmp://mediamtx:1935/live` (hostname internal Docker), serta volume `worker_storage:/app/data` yang **harus dibagi** dengan container worker-core.

### Regenerasi protobuf

```bash
python -m grpc_tools.protoc --proto_path=. --python_out=. --grpc_python_out=. \
  nedo_vision_worker/protos/*.proto
```

> ⚠️ Setelah regen, pastikan runtime `protobuf` cocok dengan gencode (repo ini terkunci ke 6.31.1 — mismatch = crash import stub).

---

## 8. Deployment & CI/CD

```mermaid
flowchart LR
    Push["git push develop/main"]
    subgraph CI["ci.yml (ubuntu)"]
        Ruff["ruff check<br/>(continue-on-error)"]
        Compile["compileall"]
        Test["pytest"]
    end
    subgraph BUILD["docker-build-and-push.yml (self-hosted macOS)"]
        Colima["build amd64<br/>Colima + Rosetta"]
        Push2["push GHCR<br/>ghcr.io/tekinfopg/*"]
        Tag["tag :latest"]
    end
    WT["🐋 Watchtower<br/>(edge, pull ~30s)"]
    Prod["🌐 iris.petrokimia-gresik.com"]

    Push --> Ruff --> Compile --> Test
    Push --> Colima --> Push2 --> Tag --> WT --> Prod
```

- **CI (`ci.yml`):** `ruff check` (belum gate keras) → `compileall` → `pytest`. Trigger di push/PR ke `develop`/`main`.
- **Image (`docker-build-and-push.yml`):** dibangun di **runner self-hosted macOS** (Colima + Rosetta untuk amd64, menghindari QEMU/keburu Actions budget habis — lihat memory CI budget), push ke **GHCR** `ghcr.io/tekinfopg/*`, tag `:latest`.
- **Rollout:** **Watchtower** di mesin edge menarik image baru ± 30 detik setelah push. Prod internal: `iris.petrokimia-gresik.com` (dashboard `/app`), gRPC `50051` direct host.

---

## 9. Observability & Known Issues

### Observability

- **Heartbeat metrik:** `SendSystemUsage` tiap 10-30s — CPU %, suhu CPU, RAM (total/used/free/%), GPU[] (usage/mem/suhu), latency gRPC (ms), dan `version` gabungan `svc/core`. Ini sumber "worker online + versi" di UI manager.
- **Health/latency:** `HealthCheckService.HealthCheck` di-ping tiap 10s untuk mengukur latency edge→pusat.
- **Log:** stdout terstruktur (`[LEVEL] message`), token selalu di-redact. Tidak ada Prometheus/Loki lokal — observability terpusat di manager (lihat `ai-vision-infrastructure`).

### Known issues (dari operasi & kode)

- **Detection-stall (pipeline "Stopped").** Ketika manager mengirim `stop` di `nedo.worker.core.action`, `WorkerManager._stop_workers()` mematikan video/pipeline/sender/dataset thread — deteksi berhenti mengalir dan pipeline tampil *Stopped* di UI. Jika perintah `stop`/`restart` datang beruntun atau `pipeline_status_code` di SQLite tidak sinkron dengan status manager, worker bisa terjebak di kondisi berhenti tanpa `start` ulang. Belum ada watchdog/alert otomatis; pemulihan biasanya lewat `restart` manual dari manager.
- **Route-lazy.** Model & sumber di-sinkron/di-unduh secara *lazy* oleh `DataSyncWorker` (interval 10s) dan worker-core memuat rute/model saat pipeline pertama kali `start`. Akibatnya deteksi pertama setelah assign bisa tertunda (unduh model + init route) — tampak seperti stall singkat di awal. Mitigasi: pre-warm storage `model/` sebelum assign.
- **Config drift saat storage-path beda.** Jika `--storage-path` agen ≠ worker-core, buffer deteksi & `.worker_core_version` tidak terbaca → upload kosong dan versi core tidak dilaporkan.

---

## 10. Documentation Index

| Audiens | Dokumen | Lokasi |
|---|---|---|
| Awam / Manajemen | Executive Summary + BPMN | README §1–2 |
| Teknis (Engineer) | Arsitektur, gRPC/RabbitMQ catalog, struktur | README §3, §5 |
| Spesialis (ML/DevOps) | Detection flow, security, TLS, tuning | README §3.3–3.9, §9 |
| DevOps | CI/CD, Docker, deploy | README §7–8, `docker-compose.local.yml` |
| Platform | Matriks hardware (Jetson/ARM/GPU) | [`PLATFORM_SUPPORT.md`](PLATFORM_SUPPORT.md) |
| Developer | Regen protobuf | [`README_DEV.md`](README_DEV.md) |

---

## 11. Contact & License

- **Tech Lead:** Yafi Anshori
- **Org GitHub:** [tekinfopg](https://github.com/tekinfopg)
- **Ekosistem:** IRIS — AI Vision Platform, PT Petrokimia Gresik

**Proprietary** — © 2026 PT Petrokimia Gresik. Penggunaan internal. Tidak untuk distribusi publik.

---

<div align="center">

**Agen edge yang menjaga mata AI tetap terbuka di lantai pabrik**

*Bagian dari ekosistem IRIS — dikelola tim Tekinfo PG*

</div>
