Skip to main content

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 gRPC RabbitMQ Cache Version License


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).

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 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

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 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

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 protobuf cocok 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 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
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)

Source distribution for iris-vision-worker 1.2.1
File Size Uploaded
iris_vision_worker-1.2.1.tar.gz 157.1 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for iris-vision-worker 1.2.1
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

1.2.1 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page