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
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-serviceadalah 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).
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.
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.
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
VisionWorkerServiceStubdan stub deteksi khusus. Channel dibuatmake_grpc_channel()— insecure by default, TLS bilaGRPC_TLS_ENABLED=true. Token dikirim dua jalur: di metadataauthorization: Bearer <token>(TokenAuthInterceptor) dan di body request (back-compat).GrpcClientBasepunya 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 dandurable; queue workerauto_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-pathyang 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.dbdanfiles/detection_image/. - worker-service → membaca & mengunggah:
DataSenderWorkermenjalankan*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_versiondari storage bersama → laporkan"svc 1.2.1 / core X.Y.Z"di fieldversionheartbeat. Jadi manager/UI tahu versi keduanya walau prosesnya terpisah.
3.6 ERD — cache SQLite lokal
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
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 menyembunyikanauthorization/bearer/tokensebelum ditulis. Configtoken/passworddi-mask (***) saat di-print. - Auth-failure hard stop: RPC
UNAUTHENTICATED/PERMISSION_DENIED→set_auth_failure_callbackmen-triggerWorkerService.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,8554di-EXPOSEdi 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
python --version # 3.8+ (3.11 disarankan, sesuai image)
ffmpeg -version # wajib (RTSP/RTMP)
docker --version # 24+ (opsional, untuk container)
Setup
# 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)
# 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
python -m grpc_tools.protoc --proto_path=. --python_out=. --grpc_python_out=. \
nedo_vision_worker/protos/*.proto
⚠️ Setelah regen, pastikan runtime
protobufcocok dengan gencode (repo ini terkunci ke 6.31.1 — mismatch = crash import stub).
8. Deployment & CI/CD
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 kedevelop/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 GHCRghcr.io/tekinfopg/*, tag:latest. - Rollout: Watchtower di mesin edge menarik image baru ± 30 detik setelah push. Prod internal:
iris.petrokimia-gresik.com(dashboard/app), gRPC50051direct host.
9. Observability & Known Issues
Observability
- Heartbeat metrik:
SendSystemUsagetiap 10-30s — CPU %, suhu CPU, RAM (total/used/free/%), GPU[] (usage/mem/suhu), latency gRPC (ms), danversiongabungansvc/core. Ini sumber "worker online + versi" di UI manager. - Health/latency:
HealthCheckService.HealthCheckdi-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 (lihatai-vision-infrastructure).
Known issues (dari operasi & kode)
- Detection-stall (pipeline "Stopped"). Ketika manager mengirim
stopdinedo.worker.core.action,WorkerManager._stop_workers()mematikan video/pipeline/sender/dataset thread — deteksi berhenti mengalir dan pipeline tampil Stopped di UI. Jika perintahstop/restartdatang beruntun ataupipeline_status_codedi SQLite tidak sinkron dengan status manager, worker bisa terjebak di kondisi berhenti tanpastartulang. Belum ada watchdog/alert otomatis; pemulihan biasanya lewatrestartmanual 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 kalistart. Akibatnya deteksi pertama setelah assign bisa tertunda (unduh model + init route) — tampak seperti stall singkat di awal. Mitigasi: pre-warm storagemodel/sebelum assign. - Config drift saat storage-path beda. Jika
--storage-pathagen ≠ worker-core, buffer deteksi &.worker_core_versiontidak 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 |
| Developer | Regen protobuf | README_DEV.md |
11. Contact & License
- Tech Lead: Yafi Anshori
- Org GitHub: tekinfopg
- Ekosistem: IRIS — AI Vision Platform, PT Petrokimia Gresik
Proprietary — © 2026 PT Petrokimia Gresik. Penggunaan internal. Tidak untuk distribusi publik.
Agen edge yang menjaga mata AI tetap terbuka di lantai pabrik
Bagian dari ekosistem IRIS — dikelola tim Tekinfo PG
Release files for iris-vision-worker 1.2.1
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| iris_vision_worker-1.2.1.tar.gz | 157.1 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| iris_vision_worker-1.2.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 359.1 kB
Release files / iris_vision_worker-1.2.1.tar.gz
| Download URL | iris_vision_worker-1.2.1.tar.gz |
|---|---|
| Size | 157.1 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
118226ad3cd008449ead9231e4fb95f5c84de3cd74f9ea263a15b7a79d767a6e
|
|
BLAKE2b-256 checksum How to use checksums |
f268b016f5bea630195175558cffa1012189599617b70121cfa533f8e0fcb00a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/6.2.0 CPython/3.14.6
|
Release files / iris_vision_worker-1.2.1-py3-none-any.whl
| Download URL | iris_vision_worker-1.2.1-py3-none-any.whl |
|---|---|
| Size | 202.0 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
72a603e7e6913b45ab0de58a78a4f5813ec05ece6dc3a6a57a0d5d2717b67db1
|
|
BLAKE2b-256 checksum How to use checksums |
3360f20240b5f77ca57f722faddc5f9ddccb100b17e2490710a73fb14b91983c
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/6.2.0 CPython/3.14.6
|