
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..
Exemple de cazuri simple în care se poate folosi o astfel de automatizare:
- 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.
- Crearea automată a seriei de miniaturi pentru fișiere grafice, adăugarea de filigrane la fotografii, alte modificări ale imaginilor.
- 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ă).
- 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 . 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 80Abonament 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:
- din panoul de control.
- Intrăm în bucket-ul pentru care vom configura webhook-urile și apăsăm pe rotița dințată:

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

Completăm câmpurile:

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:

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 , î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 . Codul sursă al acestei funcții poate fi studiat în .
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 , și parametrul -script=script.sh, atunci scriptul va fi apelat în următorul mod:
script.sh bucketA some-file-to-bucket copyTrebuie 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-backupPentru 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-ashAcces 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
;;
esacPornim serverul:
ubuntu@ubuntu-basic-1-2-10gb:~/s3-webhook$ sudo ./s3-webhook -port 80 -
scripts/scripts/s3_backup_mcs_aws.shVerificăm cum va funcționa. Prin 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.txtVom 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.txtAcum, 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.txtConț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ă . 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
