Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions

Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions
Mașina de cafea Rube Goldberg

Arhitectura bazată pe evenimente crește eficiența costurilor resurselor utilizate, deoarece acestea sunt activate doar în momentul în care sunt necesare. Există multe modalități prin care acest lucru poate fi realizat fără a crea entități cloud suplimentare ca aplicații worker. Și astăzi nu voi vorbi despre FaaS, ci despre webhooks. Voi prezenta un exemplu de procesare a evenimentelor folosind webhooks pentru stocarea obiectelor.

Câteva cuvinte despre stocarea obiectelor și webhooks. Stocarea obiectelor permite stocarea oricăror date în cloud sub formă de obiecte, accesibile prin S3 sau alt API (în funcție de implementare) prin HTTP/HTTPS. Webhooks (webhooks) sunt, în general, callback-uri personalizate prin HTTP. De obicei, acestea sunt declanșate de un eveniment, de exemplu, prin trimiterea unui cod în repository sau un comentariu publicat pe blog. Când apare un eveniment, site-ul sursă trimite o cerere HTTP la URL-ul specificat pentru webhook. Ca urmare, se poate face astfel încât evenimentele de pe un site să genereze acțiuni pe altul. În cazul în care site-ul sursă este un stoc de obiecte, modificările conținutului său reprezintă evenimente.wiki.

Exemple de cazuri simple în care se poate folosi o astfel de automatizare:

  1. Crearea copiilor tuturor obiectelor într-un alt stoc de cloud. Copiile trebuie să fie create „din mers”, la orice adăugare sau modificare a fișierelor.
  2. Crearea automată a seriei de miniaturi pentru fișiere grafice, adăugarea de filigrane la fotografii, alte modificări ale imaginilor.
  3. Notificarea sosirii documentelor noi (de exemplu, un serviciu contabil distribuit încarcă rapoarte în cloud, iar monitorizarea financiară primește notificări despre rapoartele noi, le verifică și le analizează).
  4. Cazurile puțin mai complexe implică, de exemplu, generarea unei cereri către Kubernetes, care creează un pod cu containerele necesare, îi transmite parametrii sarcinii și, după procesare, închide containerul.

Ca exemplu, vom face o variantă a sarcinii 1, când modificările din bucket-ul stocării obiectelor Mail.ru Cloud Solutions (MCS) sunt sincronizate cu stocarea obiectelor AWS prin webhooks. Într-un caz real cu sarcină mare, ar trebui să se prevadă funcționarea asincronă prin înregistrarea webhooks în coadă, dar pentru sarcina de învățare vom implementa fără aceasta.

Schema de funcționare

Protocolul de interacțiune este detaliat în ghidul webhook-urilor S3 pe MCS. Schema de funcționare include următoarele elemente:

  • Serviciul de publicare, care se află pe partea S3 a stocării și publică cereri HTTP atunci când este activat webhook-ul.
  • Serverul de primire a webhook-urilor, care ascultă solicitările serviciului de publicare prin HTTP și execută acțiunile corespunzătoare. Serverul poate fi scris în orice limbaj; în exemplul nostru, vom scrie serverul în Go.

Particularitatea implementării webhook-urilor în API-ul S3 este înregistrarea serverului de primire a webhook-urilor pe serviciul de publicare. În special, serverul de primire a webhook-urilor trebuie să confirme abonarea la mesajele serviciului de publicare (în alte implementări de webhook-uri, de obicei, confirmarea abonării nu este necesară).

Prin urmare, serverul de primire a webhook-urilor trebuie să suporte două operațiuni de bază:

  • să răspundă la cererea serviciului de publicare pentru confirmarea înregistrării,
  • să proceseze evenimentele primite.

Instalarea serverului de primire a webhook-urilor

Pentru a lansa serverul de primire a webhook-urilor, este necesar un server Linux. În acest articol, vom folosi un exemplu de instanță virtuală, pe care o implementăm pe MCS.

Vom instala software-ul necesar și vom lansa serverul de primire a webhook-urilor.

ubuntu@ubuntu-basic-1-2-10gb:~$ sudo apt-get install git
Se citesc listele de pachete... Finalizat
Se construiește arborele de dependențe
Se citesc informațiile de stare... Finalizat
Următoarele pachete au fost instalate automat și nu mai sunt necesare:
 bc dns-root-data dnsmasq-base ebtables landscape-common liblxc-common 
liblxc1 libuv1 lxcfs lxd lxd-client python3-attr python3-automat 
python3-click python3-constantly python3-hyperlink
 python3-incremental python3-pam python3-pyasn1-modules 
python3-service-identity python3-twisted python3-twisted-bin 
python3-zope.interface uidmap xdelta3
Folosiți 'sudo apt autoremove' pentru a le elimina.
Pachete sugerate:
 git-daemon-run | git-daemon-sysvinit git-doc git-el git-email git-gui 
gitk gitweb git-cvs git-mediawiki git-svn
Următoarele pachete NOU vor fi instalate:
 git
0 actualizate, 1 nou instalat, 0 de eliminat și 46 neactualizate.
Trebuie să descărcați 3915 kB de arhive.
După această operațiune, vor fi folosiți 32.3 MB de spațiu disk suplimentar.
Obține:1 http://MS1.clouds.archive.ubuntu.com/ubuntu bionic-updates/main 
amd64 git amd64 1:2.17.1-1ubuntu0.7 [3915 kB]
Descărcat 3915 kB în 1s (5639 kB/s)
Selectând pachetul git, care nu fusese selectat anterior.
(Citesc baza de date ... 53932 fișiere și directoare instalate în prezent.)
Pregătind să despachetez .../git_12.17.1-1ubuntu0.7_amd64.deb ...
Despachetând git (1:2.17.1-1ubuntu0.7) ...
Configurând git (1:2.17.1-1ubuntu0.7) ...

Clonăm folderul cu serverul de primire a webhook-urilor:

ubuntu@ubuntu-basic-1-2-10gb:~$ git clone
https://github.com/RomanenkoDenys/s3-webhook.git
Clonare în 's3-webhook'...
remote: Enumerating objects: 48, done.
remote: Counting objects: 100% (48/48), done.
remote: Compressing objects: 100% (27/27), done.
remote: Total 114 (delta 20), reused 45 (delta 18), pack-reused 66
Receiving objects: 100% (114/114), 23.77 MiB | 20.25 MiB/s, done.
Resolving deltas: 100% (49/49), done.

Vom lansa serverul:

ubuntu@ubuntu-basic-1-2-10gb:~$ cd s3-webhook/
ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ sudo ./s3-webhook -port 80

Abonament pentru serviciul de publicare

Puteți înregistra serverul dumneavoastră de primire a webhook-urilor prin API sau interfața web. Pentru simplificare, vom înregistra prin interfața web:

  1. Mergem în secțiunea bucket-uri din panoul de control.
  2. Intrăm în bucket-ul pentru care vom configura webhook-urile și apăsăm pe rotița dințată:

Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions

Trecem la fila Webhooks și apăsăm Adaugă:

Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions
Completăm câmpurile:

Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions

ID — numele webhook-ului.

Event — ce evenimente să transmitem. Am setat transmiterea tuturor evenimentelor care au loc în cadrul lucrului cu fișiere (adauga și elimină).

URL — adresa serverului de primire a webhook-urilor.

Filter prefix/suffix — un filtru care permite generarea webhook-urilor doar pentru obiecte al căror nume corespunde unor reguli specifice. De exemplu, pentru a activa webhook-ul doar pentru fișiere cu extensia .png, în Filter suffix trebuie să scriem „png”.

În prezent, doar porturile 80 și 443 sunt acceptate pentru accesarea serverului de primire a webhook-urilor.

Să apăsăm Adaugă hook și vom vedea următoarele:

Exemplu de aplicație driven-by-event bazată pe webhooks în stocarea de obiecte S3 a Mail.ru Cloud Solutions
Hook adăugat.

Serverul de primire a webhook-urilor arată în logs procesul de înregistrare a hook-ului:

ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ sudo ./s3-webhook -port 80
2020/06/15 12:01:14 [POST] cerere HTTP în curs de sosire de la 
95.163.216.92:42530
2020/06/15 12:01:14 Am primit timestamp: 2020-06-15T15:01:13+03:00 TopicArn: 
mcs5259999770|myfiles-ash|s3:ObjectCreated:*,s3:ObjectRemoved:* Token: 
E2itMqAMUVVZc51pUhFWSp13DoxezvRxkUh5P7LEuk1dEe9y URL: 
http://89.208.199.220/webhook
2020/06/15 12:01:14 Generare semnătură răspuns: 
3754ce36636f80dfd606c5254d64ecb2fd8d555c27962b70b4f759f32c76b66d

Înregistrarea s-a încheiat. În secțiunea următoare vom examina mai detaliat algoritmul de funcționare al serverului de primire a webhook-urilor.

Descrierea serverului de primire a webhook-urilor

În exemplul nostru, serverul este scris în Go. Vom analiza principiile de bază ale funcționării sale.

package main

// Generate hmac_sha256_hex
func HmacSha256hex(message string, secret string) string {
}

// Generate hmac_sha256
func HmacSha256(message string, secret string) string {
}

// Send subscription confirmation
func SubscriptionConfirmation(w http.ResponseWriter, req *http.Request, body []byte) {
}

// Send subscription confirmation
func GotRecords(w http.ResponseWriter, req *http.Request, body []byte) {
}

// Liveness probe
func Ping(w http.ResponseWriter, req *http.Request) {
    // log request
    log.Printf("[%s] cerere HTTP Ping în curs de sosire de la %sn", req.Method, req.RemoteAddr)
    fmt.Fprintf(w, "Pongn")
}

//Webhook
func Webhook(w http.ResponseWriter, req *http.Request) {
}

func main() {

    // obține argumentele din linia de comandă
    bindPort := flag.Int("port", 80, "număr între 1-65535")
    bindAddr := flag.String("address", "", "adresă IP în format punctat")
    flag.StringVar(&actionScript, "script", "", "script extern de executat")
    flag.Parse()

    http.HandleFunc("/ping", Ping)
    http.HandleFunc("/webhook", Webhook)

log.Fatal(http.ListenAndServe(*bindAddr+":"+strconv.Itoa(*bindPort), nil))
}

Să analizăm funcțiile principale:

  • Ping() — ruta care răspunde la URL/ping, o implementare simplă a probei de disponibilitate.
  • Webhook() — ruta principală, managerul URL-ului/webhook-ului:
    • confirmă înregistrarea în serviciul de publicare (trecerea la funcția SubscriptionConfirmation),
    • procesează webhook-urile primite (funcția Gotrecords).
  • Funcțiile HmacSha256 și HmacSha256hex — implementări ale algoritmilor de criptare HMAC-SHA256 și HMAC-SHA256 cu ieșire sub formă de șir de numere hexazecimale pentru calcularea semnăturii.
  • main — funcția principală, procesează parametrii liniei de comandă și înregistrează managerii URL.

Parametrii liniei de comandă acceptați de server:

  • -port — portul pe care serverul va asculta.
  • -address — adresa IP pe care serverul o va asculta.
  • -script — programul extern care este apelat pentru fiecare webhook primit.

Să analizăm mai în detaliu unele funcții:

//Webhook
func Webhook(w http.ResponseWriter, req *http.Request) {

    // Read body
    body, err := ioutil.ReadAll(req.Body)
    defer req.Body.Close()
    if err != nil {
        http.Error(w, err.Error(), 500)
        return
    }

    // log request
    log.Printf("[%s] incoming HTTP request from %sn", req.Method, req.RemoteAddr)
    // check if we got subscription confirmation request
    if strings.Contains(string(body), 
""Type":"SubscriptionConfirmation"") {
        SubscriptionConfirmation(w, req, body)
    } else {
        GotRecords(w, req, body)
    }

}

Această funcție determină ce a sosit — o cerere de confirmare a înregistrării sau un webhook. Așa cum reiese din documentation, în caz de confirmare a înregistrării, următoarea structură JSON ajunge în cererea Post:

POST http://test.com HTTP/1.1
x-amz-sns-messages-type: SubscriptionConfirmation
content-type: application/json

{
    "Timestamp":"2019-12-26T19:29:12+03:00",
    "Type":"SubscriptionConfirmation",
    "Message":"Ați ales să vă abonați la subiectul $topic. Pentru a confirma abonamentul, trebuie să răspundeți cu semnătura calculată",
    "TopicArn":"mcs2883541269|bucketA|s3:ObjectCreated:Put",
    "SignatureVersion":1,
    "Token":"RPE5UuG94rGgBH6kHXN9FUPugFxj1hs2aUQc99btJp3E49tA"
}

La această cerere trebuie să răspundeți:

content-type: application/json

{"signature":"ea3fce4bb15c6de4fec365d36bcebbc34ccddf54616d5ca12e1972f82b6d37af"}

Unde semnătura este calculată astfel:

signature = hmac_sha256(url, hmac_sha256(TopicArn, 
hmac_sha256(Timestamp, Token)))

Dacă se primește un webhook, atunci structura cererii Post arată astfel:

POST  HTTP/1.1
x-amz-sns-messages-type: SubscriptionConfirmation

{ "Records":
    [
        {
            "s3": {
                "object": {
                    "eTag":"aed563ecafb4bcc5654c597a421547b2",
                    "sequencer":1577453615,
                    "key":"some-file-to-bucket",
                    "size":100
                },
            "configurationId":"1",
            "bucket": {
                "name": "bucketA",
                "ownerIdentity": {
                    "principalId":"mcs2883541269"}
                },
                "s3SchemaVersion":"1.0"
            },
            "eventVersion":"1.0",
            "requestParameters":{
                "sourceIPAddress":"185.6.245.156"
            },
            "userIdentity": {
                "principalId":"2407013e-cbc1-415f-9102-16fb9bd6946b"
            },
            "eventName":"s3:ObjectCreated:Put",
            "awsRegion":"ru-msk",
            "eventSource":"aws:s3",
            "responseElements": {
                "x-amz-request-id":"VGJR5rtJ"
            }
        }
    ]
}

În funcție de cerere, trebuie să înțelegem cum să procesăm datele. Am ales ca indicator înregistrarea "Type":"SubscriptionConfirmation", deoarece este prezent în cererea de confirmare a abonamentului și nu este prezent în webhook. Având în vedere prezența/absența acestei înregistrări în cererea POST, execuția ulterioară a programului trece fie la funcția SubscriptionConfirmation, fie la funcția GotRecords.

Nu vom analiza în detaliu funcția SubscriptionConfirmation, aceasta este implementată conform principiilor prezentate în documentation. Codul sursă al acestei funcții poate fi studiat în repo-ul git al proiectului.

Funcția GotRecords analizează cererea primită și pentru fiecare obiect Record apelează un script extern (numele acestuia a fost transmis în parametru -script) cu următoarele parametre:

  • numele bucket-ului
  • cheia obiectului
  • acțiunea:
    • copy — dacă în cererea inițială EventName = ObjectCreated | PutObject | PutObjectCopy
    • delete — dacă în cererea inițială EventName = ObjectRemoved | DeleteObject

Astfel, dacă a fost primit un webhook cu o cerere POST, așa cum este descris mai sus, și parametrul -script=script.sh, atunci scriptul va fi apelat în următorul mod:

script.sh  bucketA some-file-to-bucket copy

Trebuie să înțelegem că acest server de primire webhook-uri nu reprezintă o soluție de producție finalizată, ci un exemplu simplificat al unei posibile implementări.

Exemplul de lucru

Vom face sincronizarea fișierelor din bucket-ul principal în MCS în bucket-ul rezervat în AWS. Bucket-ul principal se numește myfiles-ash, iar cel rezervat — myfiles-backup (configurarea bucket-ului în AWS depășește subiectul acestui articol). Prin urmare, când un fișier este plasat în bucket-ul principal, copia sa trebuie să apară în cel rezervat, iar când este șters din principal — să fie șters din rezervat.

Vom lucra cu bucket-uri folosind utilitarul awscli, care este compatibil atât cu stocarea cloud MCS, cât și cu stocarea cloud AWS.

ubuntu@ubuntu-basic-1-2-10gb:~$ sudo apt-get install awscli
Reading package lists... Done
Building dependency tree
Reading state information... Done
After this operation, 34.4 MB of additional disk space will be used.
Unpacking awscli (1.14.44-1ubuntu1) ...
Setting up awscli (1.14.44-1ubuntu1) ...

Vom configura accesul la API S3 MCS:

ubuntu@ubuntu-basic-1-2-10gb:~$ aws configure --profile mcs
AWS Access Key ID [None]: hdywEPtuuJTExxxxxxxxxxxxxx
AWS Secret Access Key [None]: hDz3SgxKwXoxxxxxxxxxxxxxxxxxx
Default region name [None]:
Default output format [None]:

Vom configura accesul la API S3 AWS:

ubuntu@ubuntu-basic-1-2-10gb:~$ aws configure --profile aws
AWS Access Key ID [None]: AKIAJXXXXXXXXXXXX
AWS Secret Access Key [None]: dfuerphOLQwu0CreP5Z8l5fuXXXXXXXXXXXXXXXX
Default region name [None]:
Default output format [None]:

Vom verifica accesul:

La AWS:

ubuntu@ubuntu-basic-1-2-10gb:~$ aws s3 ls --profile aws
2020-07-06 08:44:11 myfiles-backup

Pentru MCS, în timpul utilizării comenzii, trebuie să adăugăm —endpoint-url:

ubuntu@ubuntu-basic-1-2-10gb:~$ aws s3 ls --profile mcs --endpoint-url 
https://hb.bizmrg.com
2020-02-04 06:38:05 databasebackups-0cdaaa6402d4424e9676c75a720afa85
2020-05-27 10:08:33 myfiles-ash

Acces obținut.

Acum vom scrie un script pentru prelucrarea webhook-ului care vine, numit s3_backup_mcs_aws.sh

#!/bin/bash
# Require aws cli
# if file added — copy it to backup bucket
# if file removed — remove it from backup bucket
# Variables
ENDPOINT_MCS="https://hb.bizmrg.com"
AWSCLI_MCS=`which aws`" --endpoint-url ${ENDPOINT_MCS} --profile mcs s3"
AWSCLI_AWS=`which aws`" --profile aws s3"
BACKUP_BUCKET="myfiles-backup"

SOURCE_BUCKET="${1}"
SOURCE_FILE="${2}"
ACTION="${3}"

SOURCE="s3://${SOURCE_BUCKET}/${SOURCE_FILE}"
TARGET="s3://${BACKUP_BUCKET}/${SOURCE_FILE}"
TEMP="/tmp/${SOURCE_BUCKET}/${SOURCE_FILE}"

case ${ACTION} in
    "copy")
    ${AWSCLI_MCS} cp "${SOURCE}" "${TEMP}"
    ${AWSCLI_AWS} cp "${TEMP}" "${TARGET}"
    rm ${TEMP}
    ;;

    "delete")
    ${AWSCLI_AWS} rm ${TARGET}
    ;;

    *)
    echo "Usage: ${0} sourcebucket sourcefile copy/delete"
    exit 1
    ;;
esac

Pornim serverul:

ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ sudo ./s3-webhook -port 80 -
scripts/scripts/s3_backup_mcs_aws.sh

Verificăm cum va funcționa. Prin interfața web MCS vom adăuga fișierul test.txt în bucket-ul myfiles-ash. În log-uri, în consolă, se poate vedea că a fost trimisă o solicitare către serverul de webhookuri:

2020/07/06 09:43:08 [POST] cerere HTTP primită de la 
95.163.216.92:56612
download: s3://myfiles-ash/test.txt to ..///tmp/myfiles-ash/test.txt
upload: ..///tmp/myfiles-ash/test.txt to 
s3://myfiles-backup/test.txt

Vom verifica conținutul bucket-ului myfiles-backup în AWS:

ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ aws s3 --profile aws ls 
myfiles-backup
2020-07-06 09:43:10       1104 test.txt

Acum, prin interfața web, vom șterge fișierul din bucket-ul myfiles-ash.

Log-urile serverului:

2020/07/06 09:44:46 [POST] cerere HTTP primită de la 
95.163.216.92:58224
delete: s3://myfiles-backup/test.txt

Conținutul bucket-ului:

ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ aws s3 --profile aws ls 
myfiles-backup
ubuntu@ubuntu-basic-1-2-10gb:~$

Fișierul a fost șters, problema este rezolvată.

Concluzie și ToDo

Tot codul utilizat în acest articol se află în repository-ul meu. Acolo se află exemple de scripturi și exemple de calculare a semnăturilor pentru înregistrarea webhook-urilor.

Acest cod este doar un exemplu de cum ar putea fi utilizate webhook-urile S3 în activitățile tale. Așa cum am menționat la început, dacă se planifică utilizarea unui astfel de server în producție, este necesar, cel puțin, să rescrii serverul pentru a funcționa în mod asincron: webhook-urile primite trebuie înregistrate într-o coadă (RabbitMQ sau NATS), iar de acolo cele trebuie prelucrate de aplicații worker. În caz contrar, la un număr mare de webhook-uri primite, te poți confrunta cu lipsa resurselor serverului pentru a finaliza sarcinile. Existența coadelor permite dispersarea serverului și a worker-ilor, precum și rezolvarea problemelor în cazul repetării sarcinilor în cazul erorilor. De asemenea, ar fi bine să schimbi jurnalizarea într-una mai detaliată și mai standardizată.

Good luck!

Mai poți citi despre subiect:

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster