No Description

orborus.go 146KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817381838193820382138223823382438253826382738283829383038313832383338343835383638373838383938403841384238433844384538463847384838493850385138523853385438553856385738583859386038613862386338643865386638673868386938703871387238733874387538763877387838793880388138823883388438853886388738883889389038913892389338943895389638973898389939003901390239033904390539063907390839093910391139123913391439153916391739183919392039213922392339243925392639273928392939303931393239333934393539363937393839393940394139423943394439453946394739483949395039513952395339543955395639573958395939603961396239633964396539663967396839693970397139723973397439753976397739783979398039813982398339843985398639873988398939903991399239933994399539963997399839994000400140024003400440054006400740084009401040114012401340144015401640174018401940204021402240234024402540264027402840294030403140324033403440354036403740384039404040414042404340444045404640474048404940504051405240534054405540564057405840594060406140624063406440654066406740684069407040714072407340744075407640774078407940804081408240834084408540864087408840894090409140924093409440954096409740984099410041014102410341044105410641074108410941104111411241134114411541164117411841194120412141224123412441254126412741284129413041314132413341344135413641374138413941404141414241434144414541464147414841494150415141524153415441554156415741584159416041614162416341644165416641674168416941704171417241734174417541764177417841794180418141824183418441854186418741884189419041914192419341944195419641974198419942004201420242034204420542064207420842094210421142124213421442154216421742184219422042214222422342244225422642274228422942304231423242334234423542364237423842394240424142424243424442454246424742484249425042514252425342544255425642574258425942604261426242634264426542664267426842694270427142724273427442754276427742784279428042814282428342844285428642874288428942904291429242934294429542964297429842994300430143024303430443054306430743084309431043114312431343144315431643174318431943204321432243234324432543264327432843294330433143324333433443354336433743384339434043414342434343444345434643474348434943504351435243534354435543564357435843594360436143624363436443654366436743684369437043714372437343744375437643774378437943804381438243834384438543864387438843894390439143924393439443954396439743984399440044014402440344044405440644074408440944104411441244134414441544164417441844194420442144224423442444254426442744284429443044314432443344344435443644374438443944404441444244434444444544464447444844494450445144524453445444554456445744584459446044614462446344644465446644674468446944704471447244734474447544764477447844794480448144824483448444854486448744884489449044914492449344944495449644974498449945004501450245034504450545064507450845094510451145124513451445154516451745184519452045214522452345244525452645274528452945304531453245334534453545364537453845394540454145424543454445454546454745484549455045514552455345544555455645574558455945604561456245634564456545664567456845694570457145724573457445754576457745784579458045814582458345844585458645874588458945904591459245934594459545964597459845994600460146024603460446054606460746084609461046114612461346144615461646174618461946204621462246234624462546264627462846294630463146324633463446354636463746384639464046414642464346444645464646474648464946504651465246534654465546564657465846594660466146624663466446654666466746684669467046714672467346744675467646774678467946804681468246834684468546864687468846894690469146924693469446954696469746984699470047014702470347044705470647074708470947104711471247134714471547164717471847194720472147224723472447254726472747284729473047314732473347344735473647374738473947404741474247434744474547464747474847494750475147524753475447554756475747584759476047614762476347644765476647674768476947704771477247734774477547764777477847794780478147824783478447854786478747884789479047914792
  1. package main
  2. /*
  3. Orborus exists to listen for new jobs from Shuffle. This is to run workflows, pipelines, and other tasks.
  4. */
  5. import (
  6. "archive/zip"
  7. "bytes"
  8. "context"
  9. "encoding/json"
  10. "errors"
  11. "fmt"
  12. "io"
  13. "io/ioutil"
  14. "log"
  15. "math"
  16. "net"
  17. "net/http"
  18. "os"
  19. "os/exec"
  20. "os/signal"
  21. "path/filepath"
  22. "runtime"
  23. "strconv"
  24. "strings"
  25. "sync"
  26. "syscall"
  27. "time"
  28. "github.com/shuffle/shuffle-shared"
  29. "math/rand"
  30. //"os/signal"
  31. //"syscall"
  32. "github.com/docker/docker/api/types"
  33. "github.com/docker/docker/api/types/container"
  34. "github.com/docker/docker/api/types/filters"
  35. "github.com/docker/docker/api/types/image"
  36. "github.com/docker/docker/api/types/mount"
  37. "github.com/docker/docker/api/types/network"
  38. "github.com/docker/docker/api/types/swarm"
  39. "github.com/docker/go-connections/nat"
  40. //"github.com/docker/docker/api/types/filters"
  41. dockerclient "github.com/docker/docker/client"
  42. uuid "github.com/satori/go.uuid"
  43. //"github.com/mackerelio/go-osstat/disk"
  44. //"github.com/mackerelio/go-osstat/memory"
  45. //"github.com/shirou/gopsutil/cpu"
  46. appsv1 "k8s.io/api/apps/v1"
  47. corev1 "k8s.io/api/core/v1"
  48. rbacv1 "k8s.io/api/rbac/v1"
  49. "k8s.io/apimachinery/pkg/api/resource"
  50. metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
  51. "k8s.io/apimachinery/pkg/util/intstr"
  52. )
  53. // Starts jobs in bulk, so this could be increased or decreased based on who the user is
  54. var sleepTime = 2
  55. // Making it work on low-end machines even during busy times :)
  56. // May cause some things to run slowly
  57. var maxConcurrency = 25
  58. // Timeout if something rashes
  59. var workerTimeoutEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_TIMEOUT")
  60. var concurrencyEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY")
  61. var appSdkVersion = os.Getenv("SHUFFLE_APP_SDK_VERSION")
  62. var workerVersion = os.Getenv("SHUFFLE_WORKER_VERSION")
  63. var newWorkerImage = os.Getenv("SHUFFLE_WORKER_IMAGE")
  64. var dockerSwarmBridgeMTU = os.Getenv("SHUFFLE_SWARM_BRIDGE_DEFAULT_MTU")
  65. var dockerSwarmBridgeInterface = os.Getenv("SHUFFLE_SWARM_BRIDGE_DEFAULT_INTERFACE")
  66. var maxCPUPercent = 90
  67. // Kubernetes settings
  68. var isKubernetes = os.Getenv("IS_KUBERNETES")
  69. var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
  70. var workerServiceAccountName = os.Getenv("SHUFFLE_WORKER_SERVICE_ACCOUNT_NAME")
  71. var workerPodSecurityContext = os.Getenv("SHUFFLE_WORKER_POD_SECURITY_CONTEXT")
  72. var workerContainerSecurityContext = os.Getenv("SHUFFLE_WORKER_CONTAINER_SECURITY_CONTEXT")
  73. var appServiceAccountName = os.Getenv("SHUFFLE_APP_SERVICE_ACCOUNT_NAME")
  74. var appPodSecurityContext = os.Getenv("SHUFFLE_APP_POD_SECURITY_CONTEXT")
  75. var appContainerSecurityContext = os.Getenv("SHUFFLE_APP_CONTAINER_SECURITY_CONTEXT")
  76. var debug = os.Getenv("DEBUG") == "true"
  77. // var baseimagename = "docker.pkg.github.com/shuffle/shuffle"
  78. // var baseimagename = "ghcr.io/frikky"
  79. // var baseimagename = "shuffle/shuffle"
  80. var baseimageregistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")
  81. var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME")
  82. //var baseimagetagsuffix = os.Getenv("SHUFFLE_BASE_IMAGE_TAG_SUFFIX")
  83. // Used for cloud with auth. Onprem in certain cases too.
  84. var auth = os.Getenv("AUTH")
  85. var org = os.Getenv("ORG")
  86. // var orgId = os.Getenv("ORG_ID")
  87. var baseUrl = os.Getenv("BASE_URL")
  88. var workerServerUrl = os.Getenv("SHUFFLE_WORKER_SERVER_URL")
  89. var environment = os.Getenv("ENVIRONMENT_NAME")
  90. var dockerApiVersion = os.Getenv("DOCKER_API_VERSION")
  91. var runningMode = strings.ToLower(os.Getenv("RUNNING_MODE"))
  92. var cleanupEnv = strings.ToLower(os.Getenv("CLEANUP"))
  93. var timezone = os.Getenv("TZ")
  94. var containerName = os.Getenv("ORBORUS_CONTAINER_NAME")
  95. var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG")
  96. var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
  97. var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
  98. var memcached = os.Getenv("SHUFFLE_MEMCACHED")
  99. var queuePerMinute = os.Getenv("SHUFFLE_EXECUTION_PER_MINIUTE")
  100. var queuePerMinuteInt int
  101. // Used to download file categories. Not required since 2.1.1
  102. var pipelineApikey = ""
  103. var pipelineUrl = os.Getenv("SHUFFLE_PIPELINE_URL")
  104. var executionIds = []string{}
  105. var pipelines = []shuffle.PipelineInfo{}
  106. var namespacemade = false // For K8s
  107. var skipPipelineMount = false
  108. var tenzirDisabled = false
  109. var dockercli *dockerclient.Client
  110. var containerId string
  111. var executionCount = 0
  112. var orborusUuid = os.Getenv("SHUFFLE_ORBORUS_UUID")
  113. var imagedownloadTimeout = time.Second * 300
  114. var window = shuffle.NewTimeWindow(1 * time.Minute)
  115. func init() {
  116. var err error
  117. // dockercli, err = dockerclient.NewEnvClient()
  118. dockercli, dockerApiVersion, err = shuffle.GetDockerClient()
  119. if err != nil {
  120. log.Printf("Unable to create docker client: %s", err)
  121. }
  122. if os.Getenv("SHUFFLE_EC2_INSTANCE") == "true" {
  123. log.Printf("[INFO] Detected AWS EC2 instance. Setting up Docker Swarm with AWS optimizations.")
  124. containers, err := dockercli.ContainerList(context.Background(), container.ListOptions{})
  125. if err == nil {
  126. for _, container := range containers {
  127. if strings.Contains(container.Image, "shuffle-orborus") {
  128. if len(container.Names) != 0 {
  129. if strings.Contains(container.Names[0], "shuffle-orborus") {
  130. containerName = container.Names[0]
  131. containerName = strings.TrimPrefix(containerName, "/")
  132. os.Setenv("ORBORUS_CONTAINER_NAME", containerName)
  133. log.Printf("[DEBUG] Found orborus container name: %s", containerName)
  134. break
  135. }
  136. }
  137. }
  138. }
  139. } else {
  140. log.Printf("[ERROR] Failed to find orborus container: %s", err)
  141. }
  142. }
  143. getThisContainerId()
  144. if len(pipelineApikey) == 0 {
  145. if len(os.Getenv("SHUFFLE_AUTHORIZATION")) > 0 {
  146. log.Printf("[DEBUG] No pipeline API key found. Overriding with api key from SHUFFLE_AUTHORIZATION")
  147. pipelineApikey = os.Getenv("SHUFFLE_AUTHORIZATION")
  148. os.Setenv("SHUFFLE_PIPELINE_AUTH", pipelineApikey)
  149. }
  150. }
  151. }
  152. // form id of current running container
  153. func getThisContainerId() {
  154. fCol := ""
  155. // some adjusting based on current running mode
  156. switch runningMode {
  157. case "kubernetes":
  158. // cgroup will be like:
  159. // 11:net_cls,net_prio:/kubepods/besteffort/podf132b44d-cfcf-43f7-9906-79f58e268333/851466f8b5ed5aa0f265b1c95c6d2bafbc51a38dd5c5a1621b6e586572150009
  160. fCol = "5"
  161. log.Printf("[INFO] Running containerized in Kubernetes!")
  162. case "docker":
  163. // cgroup will be like:
  164. // 12:perf_event:/docker/0f06810364f52a2cd6e80bfba27419cb8a29758a204cd676388f4913bb366f2b
  165. fCol = "3"
  166. log.Printf("[INFO] Running containerized in Docker!")
  167. default:
  168. fCol = "3" // for backward-compatibility with production
  169. log.Printf("[WARNING] RUNNING_MODE not set - defaulting to Docker (NOT Kubernetes).")
  170. }
  171. if fCol != "" {
  172. cmd := fmt.Sprintf("cat /proc/self/cgroup | grep memory | tail -1 | cut -d/ -f%s | grep -o -E '[0-9A-z]{64}'", fCol)
  173. out, err := exec.Command("bash", "-c", cmd).Output()
  174. if err == nil {
  175. containerId = strings.TrimSpace(string(out))
  176. log.Printf("[DEBUG] Set containerId network to %s", containerId)
  177. // cgroup error. Use fallback strategy below.
  178. // https://github.com/moby/moby/issues/7015
  179. //log.Printf("Checking if %s is in %s", ".scope", string(out))
  180. if strings.Contains(string(out), ".scope") {
  181. log.Printf("[DEBUG] ContainerId contains scope. setting to empty.")
  182. containerId = ""
  183. //docker-76c537e9a4b7c7233011f5d70e6b7f2d600b6413ac58a96519b8dca7a3f7117a.scope
  184. }
  185. } else {
  186. log.Printf("[WARNING] Failed getting container ID: %s", err)
  187. }
  188. }
  189. if containerId == "" {
  190. if containerName != "" {
  191. containerId = containerName
  192. log.Printf("[INFO] Falling back to ORBORUS_CONTAINER_NAME as container ID")
  193. } else {
  194. containerId = "shuffle-orborus"
  195. log.Printf(`[WARNING] ORBORUS_CONTAINER_NAME env is not set. Falling back to default name "%s" as container ID. This may cause issues on the same server`, containerId)
  196. }
  197. }
  198. log.Printf(`[INFO] Started with containerId "%s"`, containerId)
  199. }
  200. func skipCheckInCleanup(name string) bool {
  201. return strings.HasPrefix(name, "backend") ||
  202. strings.HasPrefix(name, "shuffle-backend") ||
  203. strings.HasPrefix(name, "frontend") ||
  204. strings.HasPrefix(name, "shuffle-frontend") ||
  205. strings.HasPrefix(name, "orborus") ||
  206. strings.HasPrefix(name, "shuffle-orborus") ||
  207. strings.HasPrefix(name, "opensearch") ||
  208. strings.HasPrefix(name, "shuffle-opensearch") ||
  209. strings.HasPrefix(name, "memcached") ||
  210. strings.HasPrefix(name, "shuffle-memcached")
  211. }
  212. func cleanupExistingNodes(ctx context.Context) error {
  213. if cleanupEnv != "true" {
  214. log.Printf("[INFO] Skipping cleanup of existing workers as CLEANUP is NOT set to true. Swarm actions are being auto-discovered during executions then instead.")
  215. return nil
  216. }
  217. if isKubernetes == "true" {
  218. // Cleanup all workers created by orborus and all apps created by workers.
  219. if kubernetesNamespace == "" {
  220. kubernetesNamespace = "default"
  221. }
  222. clientset, _, err := shuffle.GetKubernetesClient()
  223. if err != nil {
  224. log.Printf("[ERROR] Error getting kubernetes client:", err)
  225. return err
  226. }
  227. // Delete all services
  228. services, err := clientset.CoreV1().Services(kubernetesNamespace).List(context.Background(), metav1.ListOptions{
  229. LabelSelector: "app.kubernetes.io/name in (shuffle-worker, shuffle-app),app.kubernetes.io/managed-by in (shuffle-orborus, shuffle-worker)",
  230. })
  231. if err != nil {
  232. log.Printf("[ERROR] Failed listing services: %s", err)
  233. return err
  234. }
  235. for _, service := range services.Items {
  236. err := clientset.CoreV1().Services(kubernetesNamespace).Delete(context.Background(), service.Name, metav1.DeleteOptions{})
  237. if err != nil {
  238. log.Printf("[ERROR] Failed deleting service %s: %s", service.Name, err)
  239. }
  240. }
  241. deployments, err := clientset.AppsV1().Deployments(kubernetesNamespace).List(context.Background(), metav1.ListOptions{
  242. LabelSelector: "app.kubernetes.io/name in (shuffle-worker, shuffle-app),app.kubernetes.io/managed-by in (shuffle-orborus, shuffle-worker)",
  243. })
  244. if err != nil {
  245. log.Printf("[ERROR] Failed listing deployments: %s", err)
  246. return err
  247. }
  248. for _, deployment := range deployments.Items {
  249. err := clientset.AppsV1().Deployments(kubernetesNamespace).Delete(context.Background(), deployment.Name, metav1.DeleteOptions{})
  250. if err != nil {
  251. log.Printf("[ERROR] Failed deleting deployment %s: %s", deployment.Name, err)
  252. }
  253. }
  254. log.Printf("[INFO] Cleaned up all services and deployments in namespace %s. Waiting 10 seconds for cleanup to reflect", kubernetesNamespace)
  255. time.Sleep(10 * time.Second)
  256. return nil
  257. }
  258. serviceListOptions := types.ServiceListOptions{}
  259. services, err := dockercli.ServiceList(
  260. context.Background(),
  261. serviceListOptions,
  262. )
  263. if err != nil {
  264. log.Printf("[DEBUG] Failed finding containers: %s", err)
  265. return err
  266. }
  267. //log.Printf("\n\nFound %d contaienrs", len(services))
  268. for _, service := range services {
  269. //portFound := false
  270. //for _, endpoint := range service.Spec.EndpointSpec.Ports {
  271. // if strings.Contains(endpoint.Name, "port") {
  272. // //portFound = true
  273. // }
  274. //}
  275. if strings.Contains(service.Spec.Annotations.Name, "opensearch") {
  276. continue
  277. }
  278. if strings.Contains(service.Spec.TaskTemplate.ContainerSpec.Image, "shuffle") {
  279. if !strings.Contains(service.Spec.TaskTemplate.ContainerSpec.Image, "shuffle-frontend") &&
  280. !strings.Contains(service.Spec.TaskTemplate.ContainerSpec.Image, "shuffle-backend") &&
  281. !strings.Contains(service.Spec.TaskTemplate.ContainerSpec.Image, "shuffle-orborus") {
  282. err = dockercli.ServiceRemove(ctx, service.ID)
  283. if err != nil {
  284. log.Printf("[DEBUG] Failed to remove service %s", service.Spec.Annotations.Name)
  285. } else {
  286. log.Printf("[DEBUG] Removed service %#v", service.Spec.TaskTemplate.ContainerSpec.Image)
  287. }
  288. }
  289. }
  290. }
  291. return nil
  292. }
  293. func deployServiceWorkers(image string) {
  294. log.Printf("[DEBUG] Validating deployment of workers as services IF swarmConfig = run (value: %#v)", swarmConfig)
  295. if swarmConfig != "run" && swarmConfig != "swarm" {
  296. log.Printf("[DEBUG] Skipping deployment of workers as services as swarmConfig is not set to run or swarm. Value: %#v", swarmConfig)
  297. return
  298. }
  299. ctx := context.Background()
  300. // Looks for and cleans up all existing items in swarm we can't re-use (Shuffle only)
  301. // frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/shuffle/shuffle-worker:nightly
  302. // Get a list of network interfaces
  303. interfaces, err := net.Interfaces()
  304. if err != nil {
  305. log.Printf("[ERROR] Failed to get network interfaces: %s", err)
  306. }
  307. mtu := 1500
  308. if len(dockerSwarmBridgeMTU) == 0 {
  309. mtu, err = strconv.Atoi(dockerSwarmBridgeMTU) // by default
  310. if err != nil {
  311. if debug {
  312. log.Printf("[DEBUG] Failed to convert the default MTU to int: %s. Using 1500 instead. Input: %s", err, dockerSwarmBridgeMTU)
  313. }
  314. mtu = 1500
  315. }
  316. }
  317. bridgeName := dockerSwarmBridgeInterface
  318. if bridgeName == "" {
  319. bridgeName = "eth0"
  320. }
  321. // Check if there is at least one interface
  322. if len(interfaces) < 2 {
  323. // this assumes that the machine should have at least 2 network
  324. // interfaces. If not, we will use the default MTU.
  325. // interface 1 is the loopback interface
  326. // interface 2 is eth0, The eth0 interface inside a
  327. // Docker container corresponds to the virtual Ethernet
  328. // interface that connects the container to the docker0
  329. log.Printf("[ERROR] Failed to get enough network interfaces")
  330. } else {
  331. // Get the preferred interface
  332. for _, iface := range interfaces {
  333. if strings.Contains(iface.Name, bridgeName) {
  334. targetInterface := iface
  335. mtu = targetInterface.MTU
  336. log.Printf("[INFO] Using MTU %d from interface %s", mtu, targetInterface.Name)
  337. break
  338. }
  339. }
  340. }
  341. // Create the network options with the specified MTU
  342. options := make(map[string]string)
  343. options["com.docker.network.driver.mtu"] = fmt.Sprintf("%d", mtu)
  344. ingressOptions := network.CreateOptions{
  345. Driver: "overlay",
  346. Attachable: false,
  347. Ingress: true,
  348. IPAM: &network.IPAM{
  349. Driver: "default",
  350. Config: []network.IPAMConfig{
  351. network.IPAMConfig{
  352. Subnet: "10.225.225.0/24",
  353. Gateway: "10.225.225.1",
  354. },
  355. },
  356. },
  357. }
  358. _, err = dockercli.NetworkCreate(
  359. ctx,
  360. "ingress",
  361. ingressOptions,
  362. )
  363. if err != nil {
  364. log.Printf("[WARNING] Ingress network may already exist: %s", err)
  365. }
  366. //docker network create --driver=overlay workers
  367. // Specific subnet?
  368. networkName := "shuffle_swarm_executions"
  369. if len(swarmNetworkName) > 0 {
  370. networkName = swarmNetworkName
  371. }
  372. networkCreateOptions := network.CreateOptions{
  373. Driver: "overlay",
  374. Options: options,
  375. Attachable: true,
  376. Ingress: false,
  377. IPAM: &network.IPAM{
  378. Driver: "default",
  379. Config: []network.IPAMConfig{
  380. network.IPAMConfig{
  381. Subnet: "10.224.224.0/24",
  382. Gateway: "10.224.224.1",
  383. },
  384. },
  385. },
  386. }
  387. _, err = dockercli.NetworkCreate(
  388. ctx,
  389. networkName,
  390. networkCreateOptions,
  391. )
  392. if err != nil {
  393. if strings.Contains(fmt.Sprintf("%s", err), "already exists") {
  394. // Try patching for attachable
  395. if debug {
  396. log.Printf("[DEBUG] Network %s already exists", networkName)
  397. }
  398. } else {
  399. log.Printf("[DEBUG] Failed to create network %s for workers: %s. This is not critical, and containers will still be added", networkName, err)
  400. }
  401. }
  402. networkID := ""
  403. // find network ID
  404. networks, err := dockercli.NetworkList(ctx, network.ListOptions{})
  405. if err == nil {
  406. for _, net := range networks {
  407. if net.Name == networkName {
  408. if net.Scope == "swarm" {
  409. log.Printf("[DEBUG] Found swarm-scoped network: %s (%s)", networkName, net.ID)
  410. networkID = net.ID
  411. } else {
  412. log.Printf("[WARNING] Network %s exists but is not swarm scoped (scope=%s)", networkName, net.Scope)
  413. }
  414. break
  415. }
  416. }
  417. }
  418. /*
  419. isMemcachedRunning, err := checkMemcached(ctx, dockercli)
  420. if err != nil {
  421. log.Printf("[ERROR] Failed checking memcached: %s", err)
  422. }
  423. if isMemcachedRunning == false {
  424. log.Printf("[ERROR] Memcached is not running. Will try to deploy it.")
  425. deployMemcached(dockercli)
  426. }
  427. ip := "shuffle-cache"
  428. if len(os.Getenv("SHUFFLE_MEMCACHED")) == 0 {
  429. os.Setenv("SHUFFLE_MEMCACHED", fmt.Sprintf("%s:11211", ip))
  430. }
  431. */
  432. if networkID == "" {
  433. log.Printf("[ERROR] Network %s does not exist", networkName)
  434. networkID = networkName
  435. }
  436. defaultNetworkAttach := false
  437. if containerId != "" {
  438. log.Printf("[DEBUG] Should connect orborus container to worker network as it's running in Docker with name %#v!", containerId)
  439. // https://pkg.go.dev/github.com/docker/docker@v20.10.12+incompatible/api/types/network#EndpointSettings
  440. networkConfig := &network.EndpointSettings{}
  441. err := dockercli.NetworkConnect(ctx, networkID, containerId, networkConfig)
  442. if err != nil {
  443. log.Printf("[ERROR] Failed connecting Orborus to docker network %s: %s", networkName, err)
  444. }
  445. if len(containerId) == 64 && baseUrl == "http://shuffle-backend:5001" {
  446. log.Printf("[WARNING] Network MAY not work due to backend being %s and container length 64. Will try to attach shuffle_shuffle network", baseUrl)
  447. defaultNetworkAttach = true
  448. }
  449. }
  450. if len(os.Getenv("DOCKER_HOST")) > 0 {
  451. log.Printf("[DEBUG] Deploying docker socket proxy to the network %s as the DOCKER_HOST variable is set", networkName)
  452. listOptions := container.ListOptions{
  453. All: true,
  454. }
  455. containers, err := dockercli.ContainerList(ctx, listOptions)
  456. if err == nil {
  457. for _, container := range containers {
  458. if strings.Contains(strings.ToLower(container.Image), "docker-socket-proxy") {
  459. networkConfig := &network.EndpointSettings{}
  460. err := dockercli.NetworkConnect(ctx, networkID, container.ID, networkConfig)
  461. if err != nil {
  462. log.Printf("[ERROR] Failed connecting Docker socket proxy to docker network %s: %s", networkName, err)
  463. } else {
  464. log.Printf("[INFO] Attached the docker socket proxy to the execution network")
  465. }
  466. break
  467. }
  468. }
  469. } else {
  470. log.Printf("[ERROR] Failed listing containers when deploying socket proxy on swarm: %s", err)
  471. }
  472. //} else {
  473. // log.Printf("[ERROR] Failed listing and finding the right image for docker socket proxy: %s", err)
  474. //}
  475. }
  476. // Running 2 by default instead of 1. Higher scale mechanisms - es
  477. replicas := uint64(1)
  478. scaleReplicas := os.Getenv("SHUFFLE_SCALE_REPLICAS")
  479. if len(scaleReplicas) > 0 {
  480. tmpInt, err := strconv.Atoi(scaleReplicas)
  481. if err != nil {
  482. log.Printf("[ERROR] %s is not a valid number for replication", scaleReplicas)
  483. } else {
  484. replicas = uint64(tmpInt)
  485. }
  486. log.Printf("[DEBUG] SHUFFLE_SCALE_REPLICAS set to value %#v. Trying to overwrite default (%d/node)", scaleReplicas, replicas)
  487. }
  488. innerContainerName := fmt.Sprintf("shuffle-workers")
  489. cnt, err := findActiveSwarmNodes()
  490. if err != nil {
  491. log.Printf("[ERROR] Failed to find active swarm nodes: %s. Defaulting to 1", err)
  492. }
  493. nodeCount := uint64(1)
  494. if cnt > 0 {
  495. nodeCount = uint64(cnt)
  496. }
  497. appReplicas := os.Getenv("SHUFFLE_APP_REPLICAS")
  498. appReplicaCnt := 2
  499. if len(appReplicas) > 0 {
  500. newCnt, err := strconv.Atoi(appReplicas)
  501. if err != nil {
  502. log.Printf("[ERROR] %s is not a valid number for SHUFFLE_APP_REPLICAS", appReplicas)
  503. } else {
  504. appReplicaCnt = newCnt
  505. }
  506. }
  507. log.Printf("[DEBUG] Found %d node(s) to replicate over. Defaulting to 1 IF we can't auto-discover them.", cnt)
  508. // FIXME: From September 2025 - This is set back to 1, as this doesn't really reflect how scale works at all. It is just confusing, and makes number larger/smaller "arbitrarily" instead of using default docker scale
  509. nodeCount = 1
  510. replicatedJobs := uint64(replicas * nodeCount)
  511. log.Printf("[DEBUG] Deploying %d container(s) for worker with swarm to each node. Service name: %s. Image: %s", replicas, innerContainerName, image)
  512. if timezone == "" {
  513. timezone = "Europe/Amsterdam"
  514. }
  515. // FIXME: May not need ingress ports. Could use internal services and DNS of swarm itself
  516. // https://github.com/moby/moby/blob/e2f740de442bac52b280bc485a3ca5b31567d938/api/types/swarm/service.go#L46
  517. serviceSpec := swarm.ServiceSpec{
  518. Annotations: swarm.Annotations{
  519. Name: innerContainerName,
  520. Labels: map[string]string{},
  521. },
  522. Mode: swarm.ServiceMode{
  523. Replicated: &swarm.ReplicatedService{
  524. Replicas: &replicatedJobs,
  525. },
  526. },
  527. Networks: []swarm.NetworkAttachmentConfig{
  528. swarm.NetworkAttachmentConfig{
  529. Target: networkID,
  530. },
  531. swarm.NetworkAttachmentConfig{
  532. Target: "ingress",
  533. },
  534. },
  535. EndpointSpec: &swarm.EndpointSpec{
  536. Mode: "vip",
  537. Ports: []swarm.PortConfig{
  538. swarm.PortConfig{
  539. Protocol: swarm.PortConfigProtocolTCP,
  540. PublishMode: swarm.PortConfigPublishModeIngress,
  541. Name: "worker-port",
  542. PublishedPort: 33333,
  543. TargetPort: 33333,
  544. },
  545. },
  546. },
  547. TaskTemplate: swarm.TaskSpec{
  548. Resources: &swarm.ResourceRequirements{
  549. Reservations: &swarm.Resources{},
  550. },
  551. LogDriver: &swarm.Driver{
  552. Name: "json-file",
  553. Options: map[string]string{
  554. "max-size": "10m",
  555. },
  556. },
  557. ContainerSpec: &swarm.ContainerSpec{
  558. Image: image,
  559. Env: []string{
  560. fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", os.Getenv("SHUFFLE_SWARM_CONFIG")),
  561. fmt.Sprintf("SHUFFLE_SWARM_NETWORK_NAME=%s", networkName),
  562. fmt.Sprintf("SHUFFLE_APP_REPLICAS=%d", appReplicaCnt),
  563. fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", os.Getenv("SHUFFLE_LOGS_DISABLED")),
  564. fmt.Sprintf("DEBUG_MEMORY=%s", os.Getenv("DEBUG_MEMORY")),
  565. fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")),
  566. fmt.Sprintf("SHUFFLE_MAX_SWARM_NODES=%s", os.Getenv("SHUFFLE_MAX_SWARM_NODES")),
  567. fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", os.Getenv("SHUFFLE_BASE_IMAGE_NAME")),
  568. fmt.Sprintf("SHUFFLE_APP_REQUEST_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_REQUEST_TIMEOUT")),
  569. },
  570. //Hosts: []string{
  571. // innerContainerName,
  572. //},
  573. },
  574. RestartPolicy: &swarm.RestartPolicy{
  575. Condition: swarm.RestartPolicyConditionOnFailure,
  576. },
  577. Placement: &swarm.Placement{
  578. Constraints: []string{},
  579. },
  580. },
  581. }
  582. if defaultNetworkAttach == true || strings.ToLower(os.Getenv("SHUFFLE_DEFAULT_NETWORK_ATTACH")) == "true" {
  583. targetName := "shuffle_shuffle"
  584. isAttachable := false
  585. networks, err := dockercli.NetworkList(ctx, network.ListOptions{})
  586. if err == nil {
  587. for _, net := range networks {
  588. if net.Name == targetName {
  589. if net.Scope == "swarm" {
  590. log.Printf("[DEBUG] Found swarm-scoped network: %s", targetName)
  591. isAttachable = true
  592. } else {
  593. log.Printf("[WARNING] Network %s exist but is not swarm scoped (scope=%s)", targetName, net.Scope)
  594. }
  595. break
  596. }
  597. }
  598. }
  599. if isAttachable {
  600. log.Printf("[DEBUG] Adding network attach for network %s to worker in swarm", targetName)
  601. serviceSpec.Networks = append(serviceSpec.Networks, swarm.NetworkAttachmentConfig{
  602. Target: targetName,
  603. })
  604. // FIXM: Remove this if deployment fails?
  605. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_SWARM_OTHER_NETWORK=%s", targetName))
  606. }
  607. }
  608. if dockerApiVersion != "" {
  609. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("DOCKER_API_VERSION=%s", dockerApiVersion))
  610. }
  611. if len(os.Getenv("SHUFFLE_SCALE_REPLICAS")) > 0 {
  612. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_SCALE_REPLICAS=%s", os.Getenv("SHUFFLE_SCALE_REPLICAS")))
  613. }
  614. if len(os.Getenv("SHUFFLE_MEMCACHED")) > 0 {
  615. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_MEMCACHED=%s", os.Getenv("SHUFFLE_MEMCACHED")))
  616. }
  617. if strings.ToLower(os.Getenv("SHUFFLE_PASS_WORKER_PROXY")) == "true" {
  618. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("HTTP_PROXY=%s", os.Getenv("HTTP_PROXY")))
  619. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("HTTPS_PROXY=%s", os.Getenv("HTTPS_PROXY")))
  620. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("NO_PROXY=%s", os.Getenv("NO_PROXY")))
  621. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("no_proxy=%s", os.Getenv("no_proxy")))
  622. }
  623. if len(workerServerUrl) > 0 {
  624. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_WORKER_SERVER_URL=%s", os.Getenv("SHUFFLE_WORKER_SERVER_URL")))
  625. }
  626. // Handles backend
  627. if len(os.Getenv("BASE_URL")) > 0 {
  628. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("BASE_URL=%s", os.Getenv("BASE_URL")))
  629. }
  630. if len(os.Getenv("SHUFFLE_CLOUDRUN_URL")) > 0 {
  631. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_CLOUDRUN_URL=%s", os.Getenv("SHUFFLE_CLOUDRUN_URL")))
  632. }
  633. if len(os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD")) > 0 {
  634. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_AUTO_IMAGE_DOWNLOAD=%s", os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD")))
  635. }
  636. if len(os.Getenv("DOCKER_HOST")) > 0 {
  637. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("DOCKER_HOST=%s", os.Getenv("DOCKER_HOST")))
  638. } else {
  639. if runtime.GOOS == "windows" {
  640. serviceSpec.TaskTemplate.ContainerSpec.Mounts = []mount.Mount{
  641. mount.Mount{
  642. Source: `\\.\pipe\docker_engine`,
  643. Target: `\\.\pipe\docker_engine`,
  644. Type: mount.TypeBind,
  645. },
  646. }
  647. } else {
  648. serviceSpec.TaskTemplate.ContainerSpec.Mounts = []mount.Mount{
  649. mount.Mount{
  650. Source: "/var/run/docker.sock",
  651. Target: "/var/run/docker.sock",
  652. Type: mount.TypeBind,
  653. },
  654. }
  655. }
  656. }
  657. // Look for SHUFFLE_VOLUME_BINDS
  658. if len(os.Getenv("SHUFFLE_VOLUME_BINDS")) > 0 {
  659. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_VOLUME_BINDS=%s", os.Getenv("SHUFFLE_VOLUME_BINDS")))
  660. }
  661. overrideHttpProxy := os.Getenv("SHUFFLE_INTERNAL_HTTP_PROXY")
  662. overrideHttpsProxy := os.Getenv("SHUFFLE_INTERNAL_HTTPS_PROXY")
  663. if len(overrideHttpProxy) > 0 {
  664. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTP_PROXY=%s", overrideHttpProxy))
  665. }
  666. if len(overrideHttpsProxy) > 0 {
  667. serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTPS_PROXY=%s", overrideHttpsProxy))
  668. }
  669. serviceOptions := types.ServiceCreateOptions{}
  670. _, err = dockercli.ServiceCreate(
  671. ctx,
  672. serviceSpec,
  673. serviceOptions,
  674. )
  675. // Force deploy if it's not disabled
  676. deployTenzirNode()
  677. if err == nil {
  678. log.Printf("[DEBUG] Successfully deployed workers with %d replica(s) on %d node(s)", replicas, cnt)
  679. // wait for service to be ready
  680. time.Sleep(time.Duration(rand.Intn(4)+1) * time.Second)
  681. //log.Printf("[DEBUG] Servicecreate request: %#v %#v", service, err)
  682. // patch service network
  683. // this is an edgecase that we noticed on docker version 29
  684. // and API version 1.44
  685. services, serr := dockercli.ServiceList(ctx, types.ServiceListOptions{})
  686. if serr == nil {
  687. for _, svc := range services {
  688. if svc.Spec.Annotations.Name == innerContainerName {
  689. log.Printf("[DEBUG] Found service %s (%s) — patching network attach", innerContainerName, svc.ID)
  690. spec := svc.Spec
  691. spec.TaskTemplate.Networks = append(spec.TaskTemplate.Networks, swarm.NetworkAttachmentConfig{
  692. Target: networkID,
  693. })
  694. _, uerr := dockercli.ServiceUpdate(ctx, svc.ID, svc.Version, spec, types.ServiceUpdateOptions{})
  695. if uerr != nil {
  696. log.Printf("[WARNING] Failed to patch service %s with network %s: %v", innerContainerName, networkID, uerr)
  697. } else {
  698. log.Printf("[INFO] Successfully attached network %s to service %s", networkID, innerContainerName)
  699. }
  700. break
  701. }
  702. }
  703. } else {
  704. log.Printf("[WARNING] Failed to list services for patching network attach: %v", serr)
  705. }
  706. } else {
  707. if !strings.Contains(fmt.Sprintf("%s", err), "Already Exists") && !strings.Contains(fmt.Sprintf("%s", err), "is already in use by service") {
  708. log.Printf("[ERROR] Failed making service: %s", err)
  709. if strings.Contains(fmt.Sprintf("%s", err), "networks scoped to the swarm can be used") {
  710. log.Printf("[WARNING] Swarm network attachment failed, retrying without shuffle_shuffle")
  711. var updatedNetworks []swarm.NetworkAttachmentConfig
  712. for _, net := range serviceSpec.Networks {
  713. if net.Target != "shuffle_shuffle" {
  714. updatedNetworks = append(updatedNetworks, net)
  715. }
  716. }
  717. serviceSpec.Networks = updatedNetworks
  718. var updatedEnv []string
  719. for _, env := range serviceSpec.TaskTemplate.ContainerSpec.Env {
  720. if !strings.HasPrefix(env, "SHUFFLE_SWARM_OTHER_NETWORK=") {
  721. updatedEnv = append(updatedEnv, env)
  722. }
  723. }
  724. serviceSpec.TaskTemplate.ContainerSpec.Env = updatedEnv
  725. serviceOptions := types.ServiceCreateOptions{}
  726. _, err = dockercli.ServiceCreate(
  727. ctx,
  728. serviceSpec,
  729. serviceOptions,
  730. )
  731. if err != nil {
  732. log.Printf("[ERROR] Failed to deploy service even without shuffle_shuffle network: %s", err)
  733. }
  734. }
  735. } else {
  736. log.Printf("[WARNING] Failed deploying workers: %s", err)
  737. if len(serviceSpec.Networks) > 1 {
  738. serviceSpec.Networks = []swarm.NetworkAttachmentConfig{
  739. swarm.NetworkAttachmentConfig{
  740. Target: "shuffle_shuffle",
  741. },
  742. }
  743. _, _ = dockercli.ServiceCreate(
  744. ctx,
  745. serviceSpec,
  746. serviceOptions,
  747. )
  748. }
  749. }
  750. }
  751. }
  752. // Deploys the worker with the current available environments
  753. // https://docs.docker.com/engine/api/sdk/examples/
  754. func buildEnvVars(envMap map[string]string) []corev1.EnvVar {
  755. var envVars []corev1.EnvVar
  756. for key, value := range envMap {
  757. envVars = append(envVars, corev1.EnvVar{Name: key, Value: value})
  758. }
  759. return envVars
  760. }
  761. func buildResourcesFromEnv() corev1.ResourceRequirements {
  762. requests := corev1.ResourceList{}
  763. limits := corev1.ResourceList{}
  764. type item struct {
  765. env string
  766. resourceName corev1.ResourceName
  767. resourceList corev1.ResourceList
  768. }
  769. items := []item{
  770. // kubernetes requests
  771. {env: "SHUFFLE_WORKER_CPU_REQUEST", resourceName: corev1.ResourceCPU, resourceList: requests},
  772. {env: "SHUFFLE_WORKER_MEMORY_REQUEST", resourceName: corev1.ResourceMemory, resourceList: requests},
  773. {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", resourceName: corev1.ResourceEphemeralStorage, resourceList: requests},
  774. // kubernetes limits
  775. {env: "SHUFFLE_WORKER_CPU_LIMIT", resourceName: corev1.ResourceCPU, resourceList: limits},
  776. {env: "SHUFFLE_WORKER_MEMORY_LIMIT", resourceName: corev1.ResourceMemory, resourceList: limits},
  777. {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", resourceName: corev1.ResourceEphemeralStorage, resourceList: limits},
  778. }
  779. for _, it := range items {
  780. if value := strings.TrimSpace(os.Getenv(it.env)); value != "" {
  781. if quantity, err := resource.ParseQuantity(value); err == nil {
  782. it.resourceList[it.resourceName] = quantity
  783. } else {
  784. log.Printf("[WARNING] Cannot parse %s=%q as resource quantity: %v", it.env, value, err)
  785. }
  786. }
  787. }
  788. rr := corev1.ResourceRequirements{}
  789. if len(requests) > 0 {
  790. rr.Requests = requests
  791. }
  792. if len(limits) > 0 {
  793. rr.Limits = limits
  794. }
  795. return rr
  796. }
  797. func handleBackendImageDownload(ctx context.Context, images string) error {
  798. // Replicate images with lowercase, as the name may be wrong
  799. // Most of the time lowercase is correct. Swapping to have that first
  800. originalImages := images
  801. images = strings.ToLower(images) + "," + originalImages
  802. // Remove the image
  803. handled := []string{}
  804. //log.Printf("[DEBUG] Removing existing image (s): %s", images)
  805. newImages := []string{}
  806. successful := []string{}
  807. for _, curimage := range strings.Split(images, ",") {
  808. curimage = strings.TrimSpace(curimage)
  809. if shuffle.ArrayContains(handled, curimage) {
  810. continue
  811. }
  812. handled = append(handled, curimage)
  813. if !strings.Contains(curimage, "/") {
  814. curimage = fmt.Sprintf("frikky/shuffle:%s", curimage)
  815. }
  816. newImages = append(newImages, curimage)
  817. // Force remove the current image to avoid cached layers
  818. // if swarmConfig == "run" || swarmConfig == "swarm" {
  819. // _, err := dockercli.ImageRemove(ctx, curimage, image.RemoveOptions{
  820. // Force: true,
  821. // PruneChildren: true,
  822. // })
  823. // if err != nil {
  824. // log.Printf("[ERROR] Failed removing image for re-download: %s", err)
  825. // } else {
  826. // log.Printf("[DEBUG] Removed image: %s", curimage)
  827. // }
  828. // } else {
  829. // //log.Printf("[DEBUG] Skipping image removal for %s as swarmConfig is not set to run or swarm. Value: %#v", curimage, swarmConfig)
  830. // }
  831. err := shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, curimage)
  832. if err != nil {
  833. //log.Printf("[ERROR] Failed downloading image: %s", err)
  834. } else {
  835. //log.Printf("[DEBUG] Downloaded image: %s", curimage)
  836. successful = append(successful, curimage)
  837. }
  838. }
  839. if len(successful) == 0 {
  840. log.Printf("[ERROR] Failed downloading image copies: %s. This means the app may not have been updated.", strings.Join(handled, ", "))
  841. } else {
  842. log.Printf("[DEBUG] Successfully downloaded image copies: %s", strings.Join(successful, ", "))
  843. }
  844. if swarmConfig == "run" || swarmConfig == "swarm" {
  845. log.Printf("[DEBUG] Should update service with new image after updating(s): %s. \n\nBETA REPLACEMENT IMPLEMENTATION: Contact support@shuffler.io for support.", strings.Join(newImages, "\n"))
  846. // 1. Download the image
  847. // 2. Find the existing service using the image
  848. // 3. Update the service with the new image in a rolling restart
  849. // Find the existing service
  850. serviceListOptions := types.ServiceListOptions{}
  851. services, err := dockercli.ServiceList(
  852. ctx,
  853. serviceListOptions,
  854. )
  855. if err != nil {
  856. log.Printf("[ERROR] Failed finding services: %s", err)
  857. } else {
  858. found := false
  859. for _, service := range services {
  860. //log.Printf("Service image: %s", service.Spec.TaskTemplate.ContainerSpec.Image)
  861. for _, image := range newImages {
  862. if !strings.Contains(service.Spec.TaskTemplate.ContainerSpec.Image, image) {
  863. continue
  864. }
  865. log.Printf("[DEBUG] Found service for image: %#v", service.Spec.Annotations.Name)
  866. // Update the service to run with the new image
  867. //docker service update --image username/imagename:latest servicename --force
  868. serviceUpdateOptions := types.ServiceUpdateOptions{}
  869. service.Spec.TaskTemplate.ForceUpdate++
  870. resp, err := dockercli.ServiceUpdate(
  871. ctx,
  872. service.ID,
  873. service.Version,
  874. service.Spec,
  875. serviceUpdateOptions,
  876. )
  877. if err != nil {
  878. log.Printf("[ERROR] Failed updating service %s with the new image %s: %s. Resp: %#v", service.Spec.Annotations.Name, image, err, resp)
  879. } else {
  880. log.Printf("[DEBUG] Updated service %s with the new image %s. Resp: %#v", service.Spec.Annotations.Name, image, resp)
  881. found = true
  882. if !strings.Contains(fmt.Sprintf("%s", resp), "error") {
  883. break
  884. } else {
  885. log.Printf("[ERROR] Failed updating service %s with the new image %s: %s. Resp: %#v", service.Spec.Annotations.Name, image, err, resp)
  886. }
  887. }
  888. }
  889. if found {
  890. break
  891. }
  892. }
  893. if !found {
  894. log.Printf("[DEBUG] Failed to find service to update for service %s", newImages)
  895. }
  896. }
  897. }
  898. return nil
  899. }
  900. func fixk8sRoles() {
  901. clientset, _, err := shuffle.GetKubernetesClient()
  902. if err != nil {
  903. log.Printf("[ERROR] Error getting kubernetes client: %s", err)
  904. os.Exit(1)
  905. }
  906. kubernetesNamespace := "default"
  907. // Check if namespace exist as variable. If so, make it
  908. if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 {
  909. kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
  910. }
  911. // fix roles
  912. // check if "service-creator" role is assigned to the service account "default"
  913. // roleBindingNames := []string{"service-creator-binding", "pod-creator-binding", "deployment-creator-binding"}
  914. serviceAccountName := "default"
  915. roleBindingName := "creator-all"
  916. resourceTypes := []string{"services", "pods", "deployments"}
  917. // Check if the RoleBinding exists
  918. roleBinding, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Get(context.TODO(), roleBindingName, metav1.GetOptions{})
  919. if err != nil {
  920. log.Printf("[WARNING] Failed to get RoleBinding %s: %s", roleBindingName, err)
  921. // create role and rolebinding
  922. role := &rbacv1.Role{
  923. ObjectMeta: metav1.ObjectMeta{
  924. Name: roleBindingName,
  925. },
  926. Rules: []rbacv1.PolicyRule{
  927. {
  928. APIGroups: []string{"", "apps"},
  929. Resources: resourceTypes,
  930. Verbs: []string{"create", "list"},
  931. },
  932. },
  933. }
  934. ctx := context.TODO()
  935. _, err := clientset.RbacV1().Roles(kubernetesNamespace).Create(ctx, role, metav1.CreateOptions{})
  936. if err != nil {
  937. log.Printf("[ERROR] Failed to create Role %s: %s", roleBindingName, err)
  938. if !strings.Contains(fmt.Sprintf("%s", err), "already exists") {
  939. log.Printf("[INFO] role %s already exists", roleBindingName)
  940. }
  941. }
  942. roleBinding := &rbacv1.RoleBinding{
  943. ObjectMeta: metav1.ObjectMeta{
  944. Name: roleBindingName,
  945. },
  946. Subjects: []rbacv1.Subject{
  947. {
  948. Kind: "ServiceAccount",
  949. Name: serviceAccountName,
  950. Namespace: kubernetesNamespace,
  951. },
  952. },
  953. RoleRef: rbacv1.RoleRef{
  954. Kind: "Role",
  955. Name: roleBindingName,
  956. },
  957. }
  958. _, err = clientset.RbacV1().RoleBindings(kubernetesNamespace).Create(ctx, roleBinding, metav1.CreateOptions{})
  959. if err != nil {
  960. log.Printf("[ERROR] Failed to create RoleBinding %s: %s", roleBindingName, err)
  961. if strings.Contains(fmt.Sprintf("%s", err), "already exists") {
  962. log.Printf("[INFO] rolebinding %s already exists", roleBindingName)
  963. }
  964. }
  965. log.Printf("[INFO] Created Role %s and RoleBinding %s", roleBindingName, roleBindingName)
  966. } else {
  967. log.Printf("[INFO] RoleBinding %s exists", roleBindingName)
  968. }
  969. // Check if the RoleBinding is assigned to the service account
  970. var found bool
  971. for _, subject := range roleBinding.Subjects {
  972. if subject.Kind == "ServiceAccount" && subject.Name == serviceAccountName {
  973. found = true
  974. break
  975. }
  976. }
  977. if !found {
  978. log.Printf("[WARNING] Service account %s is not assigned to RoleBinding %s\n", serviceAccountName, roleBindingName)
  979. // assign the service account to the rolebinding
  980. roleBinding.Subjects = append(roleBinding.Subjects, rbacv1.Subject{
  981. Kind: "ServiceAccount",
  982. Name: serviceAccountName,
  983. Namespace: kubernetesNamespace,
  984. })
  985. ctx := context.TODO()
  986. _, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Update(ctx, roleBinding, metav1.UpdateOptions{})
  987. if err != nil {
  988. log.Printf("[ERROR](ns - %s) Failed to update RoleBinding %s: %s", kubernetesNamespace, roleBindingName, err)
  989. if !strings.Contains(fmt.Sprintf("%s", err), "already exists") {
  990. log.Printf("[INFO] rolebinding %s already exists", roleBindingName)
  991. }
  992. }
  993. }
  994. }
  995. // TODO: Check if deployment or service already exist by labels and only create if not already exists
  996. func deployK8sWorker(image string, identifier string, env []string) error {
  997. env = append(env, fmt.Sprintf("IS_KUBERNETES=true"))
  998. env = append(env, fmt.Sprintf("KUBERNETES_NAMESPACE=%s", os.Getenv("KUBERNETES_NAMESPACE")))
  999. // app resource env
  1000. for _, k := range []string{
  1001. "SHUFFLE_APP_CPU_REQUEST",
  1002. "SHUFFLE_APP_MEMORY_REQUEST",
  1003. "SHUFFLE_APP_EPHEMERAL_STORAGE_REQUEST",
  1004. "SHUFFLE_APP_CPU_LIMIT",
  1005. "SHUFFLE_APP_MEMORY_LIMIT",
  1006. "SHUFFLE_APP_EPHEMERAL_STORAGE_LIMIT",
  1007. } {
  1008. if v := os.Getenv(k); v != "" {
  1009. env = append(env, fmt.Sprintf("%s=%s", k, v))
  1010. }
  1011. }
  1012. if len(os.Getenv("KUBERNETES_SERVICE_HOST")) > 0 {
  1013. env = append(env, fmt.Sprintf("KUBERNETES_SERVICE_HOST=%s", os.Getenv("KUBERNETES_SERVICE_HOST")))
  1014. }
  1015. if len(os.Getenv("SHUFFLE_MEMCACHED")) > 0 {
  1016. env = append(env, fmt.Sprintf("SHUFFLE_MEMCACHED=%s", os.Getenv("SHUFFLE_MEMCACHED")))
  1017. }
  1018. if len(os.Getenv("KUBERNETES_SERVICE_PORT")) > 0 {
  1019. env = append(env, fmt.Sprintf("KUBERNETES_SERVICE_PORT=%s", os.Getenv("KUBERNETES_SERVICE_PORT")))
  1020. }
  1021. if len(os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")) > 0 {
  1022. env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_REGISTRY=%s", os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")))
  1023. }
  1024. if len(os.Getenv("SHUFFLE_BASE_IMAGE_NAME")) > 0 {
  1025. env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", os.Getenv("SHUFFLE_BASE_IMAGE_NAME")))
  1026. } else {
  1027. log.Printf("[INFO] SHUFFLE_BASE_IMAGE_NAME is not set. Defaulting to %s", baseimagename)
  1028. env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", baseimagename))
  1029. }
  1030. if len(os.Getenv("REGISTRY_URL")) > 0 {
  1031. env = append(env, fmt.Sprintf("REGISTRY_URL=%s", os.Getenv("REGISTRY_URL")))
  1032. }
  1033. if len(os.Getenv("SHUFFLE_USE_GHCR_OVERRIDE_FOR_AUTODEPLOY")) > 0 {
  1034. env = append(env, fmt.Sprintf("SHUFFLE_USE_GHCR_OVERRIDE_FOR_AUTODEPLOY=%s", os.Getenv("SHUFFLE_USE_GHCR_OVERRIDE_FOR_AUTODEPLOY")))
  1035. }
  1036. if len(os.Getenv("SHUFFLE_APP_EXPOSED_PORT")) > 0 {
  1037. env = append(env, fmt.Sprintf("SHUFFLE_APP_EXPOSED_PORT=%s", os.Getenv("SHUFFLE_APP_EXPOSED_PORT")))
  1038. }
  1039. if len(appServiceAccountName) > 0 {
  1040. env = append(env, fmt.Sprintf("SHUFFLE_APP_SERVICE_ACCOUNT_NAME=%s", appServiceAccountName))
  1041. }
  1042. if len(appPodSecurityContext) > 0 {
  1043. env = append(env, fmt.Sprintf("SHUFFLE_APP_POD_SECURITY_CONTEXT=%s", appPodSecurityContext))
  1044. }
  1045. if len(appContainerSecurityContext) > 0 {
  1046. env = append(env, fmt.Sprintf("SHUFFLE_APP_CONTAINER_SECURITY_CONTEXT=%s", appContainerSecurityContext))
  1047. }
  1048. if len(os.Getenv("SHUFFLE_APP_MOUNT_TMP_VOLUME")) > 0 {
  1049. env = append(env, fmt.Sprintf("SHUFFLE_APP_MOUNT_TMP_VOLUME=%s", os.Getenv("SHUFFLE_APP_MOUNT_TMP_VOLUME")))
  1050. }
  1051. if len(os.Getenv("SHUFFLE_LOGS_DISABLED")) > 0 {
  1052. env = append(env, fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", os.Getenv("SHUFFLE_LOGS_DISABLED")))
  1053. }
  1054. if len(os.Getenv("SHUFFLE_APP_REPLICAS")) > 0 {
  1055. env = append(env, fmt.Sprintf("SHUFFLE_APP_REPLICAS=%s", os.Getenv("SHUFFLE_APP_REPLICAS")))
  1056. }
  1057. clientset, _, err := shuffle.GetKubernetesClient()
  1058. if err != nil {
  1059. log.Printf("[ERROR] Error getting kubernetes client:", err)
  1060. return err
  1061. }
  1062. //env = append(env, fmt.Sprintf("KUBERNETES_CONFIG=%s", config.String()))
  1063. // Check if namespace exist as variable. If so, make it
  1064. if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 && !namespacemade {
  1065. kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
  1066. // Make the namespace
  1067. namespace := &corev1.Namespace{
  1068. ObjectMeta: metav1.ObjectMeta{
  1069. Name: os.Getenv("KUBERNETES_NAMESPACE"),
  1070. },
  1071. }
  1072. _, err := clientset.CoreV1().Namespaces().Create(context.Background(), namespace, metav1.CreateOptions{})
  1073. if err != nil {
  1074. if !strings.Contains(strings.ToLower(fmt.Sprintf("%s", err)), "already exists") {
  1075. log.Printf("[ERROR] Failed creating Kubernetes namespace: %s", err)
  1076. } else {
  1077. namespacemade = true
  1078. }
  1079. } else {
  1080. namespacemade = true
  1081. }
  1082. }
  1083. // Required format:
  1084. // url/org/repo/appname:tag
  1085. // url/org/repo/appname:tag
  1086. //env = append(env, fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", swarmConfig))
  1087. env = append(env, fmt.Sprintf("BASE_URL=%s", baseUrl))
  1088. env = append(env, fmt.Sprintf("SHUFFLE_SWARM_CONFIG=run"))
  1089. env = append(env, fmt.Sprintf("WORKER_HOSTNAME=%s", "shuffle-workers"))
  1090. if len(kubernetesNamespace) == 0 {
  1091. foundNamespace, err := shuffle.GetKubernetesNamespace()
  1092. if err != nil {
  1093. //log.Printf("[ERROR] Failed getting Kubernetes namespace: %s", err)
  1094. }
  1095. if len(foundNamespace) > 0 {
  1096. kubernetesNamespace = foundNamespace
  1097. os.Setenv("KUBERNETES_NAMESPACE", kubernetesNamespace)
  1098. }
  1099. }
  1100. if len(kubernetesNamespace) == 0 {
  1101. kubernetesNamespace = "default"
  1102. }
  1103. kubernetesImage := os.Getenv("SHUFFLE_KUBERNETES_WORKER")
  1104. if len(kubernetesImage) == 0 {
  1105. kubernetesImage = image
  1106. }
  1107. log.Printf("[DEBUG] Using Kubernetes worker image '%s'", kubernetesImage)
  1108. // image = "shuffle-worker:v1" //hard coded image name to test locally
  1109. envMap := make(map[string]string)
  1110. for _, envStr := range env {
  1111. parts := strings.SplitN(envStr, "=", 2)
  1112. if len(parts) == 2 {
  1113. envMap[parts[0]] = parts[1]
  1114. }
  1115. }
  1116. labels := map[string]string{
  1117. // Well-known Kubernetes labels
  1118. "app.kubernetes.io/name": "shuffle-worker",
  1119. "app.kubernetes.io/instance": identifier,
  1120. "app.kubernetes.io/part-of": "shuffle",
  1121. "app.kubernetes.io/managed-by": "shuffle-orborus",
  1122. // Keep legacy labels for backward compatibility
  1123. "container": "shuffle-worker",
  1124. }
  1125. matchLabels := map[string]string{
  1126. "app.kubernetes.io/name": "shuffle-worker",
  1127. "app.kubernetes.io/instance": identifier,
  1128. }
  1129. // Parse security contexts from env
  1130. var podSecurityContext *corev1.PodSecurityContext
  1131. var containerSecurityContext *corev1.SecurityContext
  1132. if len(workerPodSecurityContext) > 0 {
  1133. podSecurityContext = &corev1.PodSecurityContext{}
  1134. err = json.Unmarshal([]byte(workerPodSecurityContext), podSecurityContext)
  1135. if err != nil {
  1136. log.Printf("[ERROR] Failed to unmarshal worker pod security context: %v", err)
  1137. return fmt.Errorf("failed to unmarshal worker pod security context: %v", err)
  1138. }
  1139. }
  1140. if len(workerContainerSecurityContext) > 0 {
  1141. containerSecurityContext = &corev1.SecurityContext{}
  1142. err = json.Unmarshal([]byte(workerContainerSecurityContext), containerSecurityContext)
  1143. if err != nil {
  1144. log.Printf("[ERROR] Failed to unmarshal worker container security context: %v", err)
  1145. return fmt.Errorf("failed to unmarshal worker container security context: %v", err)
  1146. }
  1147. }
  1148. containerAttachment := corev1.Container{
  1149. Name: identifier,
  1150. Image: kubernetesImage,
  1151. Env: buildEnvVars(envMap),
  1152. SecurityContext: containerSecurityContext,
  1153. Resources: buildResourcesFromEnv(),
  1154. //ImagePullPolicy: "Never",
  1155. ImagePullPolicy: corev1.PullIfNotPresent,
  1156. }
  1157. if len(os.Getenv("REGISTRY_URL")) > 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_NAME")) > 0 {
  1158. log.Printf("[INFO] Setting image pull policy to Always as private registry is used.")
  1159. containerAttachment.ImagePullPolicy = corev1.PullAlways
  1160. }
  1161. podname := shuffle.GetPodName()
  1162. ctx := context.Background()
  1163. if len(podname) > 0 {
  1164. _, err := shuffle.GetCurrentPodNetworkConfig(ctx, clientset, kubernetesNamespace, podname)
  1165. if err != nil {
  1166. log.Printf("[ERROR] Failed getting current pod network: %s", err)
  1167. } else {
  1168. log.Printf("[DEBUG] Current pod found!")
  1169. // currentPodStatus = k8s.io/api/core/v1.PodStatus
  1170. }
  1171. }
  1172. // While testing:
  1173. // kubectl delete pods --all --all-namespaces; kubectl delete services --all --all-namespaces
  1174. // pod := &corev1.Pod{
  1175. // ObjectMeta: metav1.ObjectMeta{
  1176. // Name: identifier,
  1177. // Labels: containerLabels,
  1178. // },
  1179. // Spec: corev1.PodSpec{
  1180. // RestartPolicy: "Never",
  1181. // // DNSPolicy: "Default",
  1182. // DNSPolicy: corev1.DNSClusterFirst,
  1183. // // NodeSelector: map[string]string{
  1184. // // "node": "master",
  1185. // // },
  1186. // Containers: []corev1.Container{
  1187. // containerAttachment,
  1188. // },
  1189. // },
  1190. // }
  1191. // // Check if running on ARM or x86 to download the correct image
  1192. // // Get current pod's network so we can make the pod in it
  1193. // _, err = clientset.CoreV1().Pods(kubernetesNamespace).List(context.Background(), metav1.ListOptions{})
  1194. // if err != nil {
  1195. // log.Printf("[ERROR] Failed listing pods: %s", err)
  1196. // }
  1197. // createdPod, err := clientset.CoreV1().Pods(kubernetesNamespace).Create(context.Background(), pod, metav1.CreateOptions{})
  1198. // if err != nil {
  1199. // //log.Printf("[ERROR] Failed creating pod: %v", err)
  1200. // return err
  1201. // }
  1202. // log.Printf("[INFO] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace)
  1203. // // kubectl expose pod shuffle-workers --type=LoadBalancer --port=33333
  1204. // service := &corev1.Service{
  1205. // ObjectMeta: metav1.ObjectMeta{
  1206. // Name: identifier,
  1207. // },
  1208. // Spec: corev1.ServiceSpec{
  1209. // Selector: map[string]string{
  1210. // "container": "shuffle-workers",
  1211. // },
  1212. // Ports: []corev1.ServicePort{
  1213. // {
  1214. // Protocol: "TCP",
  1215. // Port: 33333,
  1216. // TargetPort: intstr.FromInt(33333),
  1217. // },
  1218. // },
  1219. // Type: corev1.ServiceTypeLoadBalancer,
  1220. // },
  1221. // }
  1222. // _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.TODO(), service, metav1.CreateOptions{})
  1223. // if err != nil {
  1224. // log.Printf("[ERROR] Failed creating service: %v", err)
  1225. // return err
  1226. // }
  1227. replicaNumberStr := os.Getenv("SHUFFLE_SCALE_REPLICAS")
  1228. replicaNumber := 1
  1229. if len(replicaNumberStr) > 0 {
  1230. tmpInt, err := strconv.Atoi(replicaNumberStr)
  1231. if err != nil {
  1232. log.Printf("[ERROR] %s is not a valid number for replication", replicaNumberStr)
  1233. } else {
  1234. replicaNumber = tmpInt
  1235. }
  1236. }
  1237. existing, err := clientset.AppsV1().Deployments(kubernetesNamespace).List(ctx, metav1.ListOptions{
  1238. LabelSelector: "app.kubernetes.io/name=shuffle-worker",
  1239. })
  1240. if err != nil {
  1241. log.Printf("[ERROR] Failed listing existing deployments: %v", err)
  1242. }
  1243. if len(existing.Items) > 0 {
  1244. log.Printf("[INFO] Found existing deployments, skipping creation")
  1245. return nil
  1246. }
  1247. replicaNumberInt32 := int32(replicaNumber)
  1248. // worker makes authenticated requests to the k8s api to create app deployments.
  1249. // Therefore, it needs to have access to the service account token.
  1250. automountServiceAccountToken := true
  1251. deployment := &appsv1.Deployment{
  1252. ObjectMeta: metav1.ObjectMeta{
  1253. Name: identifier,
  1254. Labels: labels,
  1255. },
  1256. Spec: appsv1.DeploymentSpec{
  1257. Replicas: &replicaNumberInt32,
  1258. Selector: &metav1.LabelSelector{
  1259. MatchLabels: matchLabels,
  1260. },
  1261. Template: corev1.PodTemplateSpec{
  1262. ObjectMeta: metav1.ObjectMeta{
  1263. Labels: labels,
  1264. },
  1265. Spec: corev1.PodSpec{
  1266. Containers: []corev1.Container{
  1267. containerAttachment,
  1268. },
  1269. DNSPolicy: corev1.DNSClusterFirst,
  1270. ServiceAccountName: workerServiceAccountName,
  1271. AutomountServiceAccountToken: &automountServiceAccountToken,
  1272. SecurityContext: podSecurityContext,
  1273. },
  1274. },
  1275. },
  1276. }
  1277. _, err = clientset.AppsV1().Deployments(kubernetesNamespace).Create(context.Background(), deployment, metav1.CreateOptions{})
  1278. if err != nil {
  1279. log.Printf("[ERROR] Failed creating deployment: %v", err)
  1280. return err
  1281. }
  1282. svcAppProtocol := "http"
  1283. service := &corev1.Service{
  1284. ObjectMeta: metav1.ObjectMeta{
  1285. Name: identifier,
  1286. Labels: labels,
  1287. },
  1288. Spec: corev1.ServiceSpec{
  1289. Selector: matchLabels,
  1290. Ports: []corev1.ServicePort{
  1291. {
  1292. Protocol: "TCP",
  1293. AppProtocol: &svcAppProtocol,
  1294. Port: 33333,
  1295. TargetPort: intstr.FromInt(33333),
  1296. },
  1297. },
  1298. Type: corev1.ServiceTypeClusterIP,
  1299. },
  1300. }
  1301. _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.Background(), service, metav1.CreateOptions{})
  1302. if err != nil {
  1303. log.Printf("[ERROR] Failed creating service: %v", err)
  1304. return err
  1305. }
  1306. return nil
  1307. }
  1308. func deployWorker(image string, identifier string, env []string, executionRequest shuffle.ExecutionRequest) error {
  1309. if len(os.Getenv("REGISTRY_URL")) > 0 && os.Getenv("REGISTRY_URL") != "" {
  1310. env = append(env, fmt.Sprintf("REGISTRY_URL=%s", os.Getenv("REGISTRY_URL")))
  1311. }
  1312. if swarmConfig == "run" || swarmConfig == "swarm" || isKubernetes == "true" {
  1313. // FIXME: Should we handle replies properly?
  1314. // In certain cases, a workflow may e.g. be aborted already. If it's aborted, that returns
  1315. // a 401 from the worker, which returns an error here
  1316. go sendWorkerRequest(executionRequest, image, env)
  1317. return nil
  1318. }
  1319. // Binds is the actual "-v" volume.
  1320. // Max 20% CPU every second
  1321. //CPUQuota: 25000,
  1322. //CPUPeriod: 100000,
  1323. //CPUShares: 256,
  1324. hostConfig := &container.HostConfig{
  1325. LogConfig: container.LogConfig{
  1326. Type: "json-file",
  1327. Config: map[string]string{
  1328. "max-size": "10m",
  1329. },
  1330. },
  1331. Resources: container.Resources{},
  1332. }
  1333. // This is just to test the mounting locally so
  1334. // I can control from what source I'm mounting
  1335. // the certs to. Default behaviour is:
  1336. // /certs:/certs.
  1337. certPath := "/certs"
  1338. if os.Getenv("SHUFFLE_CERT_PATH") != "" {
  1339. certPath = os.Getenv("SHUFFLE_CERT_PATH")
  1340. }
  1341. _, err := os.ReadDir(certPath)
  1342. if certPath != "" && err == nil {
  1343. certVol := mount.Mount{
  1344. Type: mount.TypeBind,
  1345. Source: certPath,
  1346. Target: "/certs",
  1347. }
  1348. hostConfig.Mounts = append(hostConfig.Mounts, certVol)
  1349. }
  1350. if len(os.Getenv("DOCKER_HOST")) == 0 {
  1351. if runtime.GOOS == "windows" {
  1352. hostConfig.Binds = []string{`\\.\pipe\docker_engine:\\.\pipe\docker_engine`}
  1353. } else {
  1354. hostConfig.Binds = []string{"/var/run/docker.sock:/var/run/docker.sock:rw"}
  1355. }
  1356. }
  1357. //var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG")
  1358. parsedUuid := uuid.NewV4()
  1359. config := &container.Config{
  1360. Image: image,
  1361. Env: env,
  1362. }
  1363. if isKubernetes != "true" {
  1364. hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId))
  1365. if strings.ToLower(cleanupEnv) == "true" {
  1366. hostConfig.AutoRemove = true
  1367. }
  1368. }
  1369. //log.Printf("[INFO] Identifier: %s", identifier)
  1370. cont, err := dockercli.ContainerCreate(
  1371. context.Background(),
  1372. config,
  1373. hostConfig,
  1374. nil,
  1375. nil,
  1376. identifier,
  1377. )
  1378. if err != nil {
  1379. if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") {
  1380. identifier = fmt.Sprintf("%s-%s", identifier, parsedUuid)
  1381. //log.Printf("[INFO] 2 - Identifier: %s", identifier)
  1382. cont, err = dockercli.ContainerCreate(
  1383. context.Background(),
  1384. config,
  1385. hostConfig,
  1386. nil,
  1387. nil,
  1388. identifier,
  1389. )
  1390. if err != nil {
  1391. log.Printf("[ERROR][%s] Container create error(2): %s", executionRequest.ExecutionId, err)
  1392. return err
  1393. }
  1394. } else {
  1395. log.Printf("[ERROR][%s] Container create error: %s", executionRequest.ExecutionId, err)
  1396. return err
  1397. }
  1398. }
  1399. // FIXME: Verbosity for testing
  1400. //log.Printf("WORKER STARTING WITH ENV: %#v", env)
  1401. ctx := context.Background()
  1402. containerStartOptions := container.StartOptions{}
  1403. err = dockercli.ContainerStart(ctx, cont.ID, containerStartOptions)
  1404. if err != nil {
  1405. // Trying to recreate and start WITHOUT network if it's possible. No extended checks. Old execution system (<0.9.30)
  1406. if strings.Contains(fmt.Sprintf("%s", err), "cannot join network") || strings.Contains(fmt.Sprintf("%s", err), "No such container") {
  1407. hostConfig.NetworkMode = ""
  1408. //container.NetworkMode(fmt.Sprintf("container:%s", containerId))
  1409. cont, err = dockercli.ContainerCreate(
  1410. context.Background(),
  1411. config,
  1412. hostConfig,
  1413. nil,
  1414. nil,
  1415. identifier+"-2",
  1416. )
  1417. if err != nil {
  1418. log.Printf("[ERROR][%s] Failed to CREATE container (2): %s", executionRequest.ExecutionId, err)
  1419. }
  1420. err = dockercli.ContainerStart(context.Background(), cont.ID, containerStartOptions)
  1421. if err != nil {
  1422. log.Printf("[ERROR][%s] Failed to start container (2): %s", executionRequest.ExecutionId, err)
  1423. }
  1424. } else {
  1425. log.Printf("[ERROR][%s] Failed initial container start. Quitting as this is NOT a simple network issue. Err: %s", executionRequest.ExecutionId, err)
  1426. }
  1427. if err != nil {
  1428. log.Printf("[ERROR][%s] Failed to start worker container in environment '%s': %s", executionRequest.ExecutionId, environment, err)
  1429. return err
  1430. } else {
  1431. log.Printf("[INFO][%s] Worker Container created (2). Runtime Location '%s': docker logs -f %s", executionRequest.ExecutionId, environment, cont.ID)
  1432. }
  1433. stats, err := dockercli.ContainerInspect(ctx, cont.ID)
  1434. if err != nil {
  1435. log.Printf("[WARNING][%s] Failed checking worker '%s': %s", executionRequest.ExecutionId, cont.ID, err)
  1436. return nil
  1437. }
  1438. containerStatus := stats.ContainerJSONBase.State.Status
  1439. if containerStatus != "running" {
  1440. log.Printf("[ERROR][%s] Status of %s is %s. Should be running. Contact support@shuffler.io if this persists.", executionRequest.ExecutionId, cont.ID, containerStatus)
  1441. }
  1442. /*
  1443. err = stopWorker(containerName)
  1444. if err != nil {
  1445. log.Printf("Failed stopping worker %s", execution.ExecutionId)
  1446. return nil
  1447. }
  1448. err = deployWorker(dockercli, workerImage, containerName, env)
  1449. if err != nil {
  1450. log.Printf("Failed executing worker %s in state %s", execution.ExecutionId, containerStatus)
  1451. return nil
  1452. }
  1453. }
  1454. */
  1455. } else {
  1456. log.Printf("[INFO][%s] New Worker created. Environment %s: docker logs %s", executionRequest.ExecutionId, environment, cont.ID)
  1457. }
  1458. return nil
  1459. }
  1460. func stopWorker(containername string) error {
  1461. ctx := context.Background()
  1462. // containers, err := cli.ContainerList(ctx, types.ContainerListOptions{
  1463. // All: true,
  1464. // })
  1465. //if err := dockercli.ContainerStop(ctx, containername, nil); err != nil {
  1466. var options container.StopOptions
  1467. if err := dockercli.ContainerStop(ctx, containername, options); err != nil {
  1468. log.Printf("[ERROR] Unable to stop container %s - running removal anyway, just in case: %s", containername, err)
  1469. }
  1470. removeOptions := container.RemoveOptions{
  1471. RemoveVolumes: true,
  1472. Force: true,
  1473. }
  1474. if err := dockercli.ContainerRemove(ctx, containername, removeOptions); err != nil {
  1475. log.Printf("[ERROR] Unable to remove container: %s", err)
  1476. }
  1477. return nil
  1478. }
  1479. func initializeImages() {
  1480. ctx := context.Background()
  1481. if appSdkVersion == "" {
  1482. appSdkVersion = "latest"
  1483. log.Printf("[INFO] SHUFFLE_APP_SDK_VERSION not defined. Defaulting to %#v", appSdkVersion)
  1484. }
  1485. if workerVersion == "" {
  1486. workerVersion = "latest"
  1487. log.Printf("[INFO] SHUFFLE_WORKER_VERSION not defined. Defaulting to %#v", workerVersion)
  1488. }
  1489. if baseimageregistry == "" {
  1490. //baseimageregistry = "ghcr.io" // Github
  1491. baseimageregistry = "docker.io" // Dockerhub
  1492. if len(os.Getenv("REGISTRY_URL")) > 0 {
  1493. baseimageregistry = os.Getenv("REGISTRY_URL")
  1494. } else {
  1495. // os.Setenv("REGISTRY_URL", baseimageregistry)
  1496. }
  1497. os.Setenv("SHUFFLE_BASE_IMAGE_REGISTRY", baseimageregistry)
  1498. log.Printf("[INFO] Setting baseimageregistry to %#v", baseimageregistry)
  1499. }
  1500. if baseimagename == "" {
  1501. // FIXME: This is probably the problem for image names tbh
  1502. //baseimagename = "shuffle" // Github (ghcr.io)
  1503. baseimagename = "frikky/shuffle" // Dockerhub
  1504. os.Setenv("SHUFFLE_BASE_IMAGE_NAME", baseimagename)
  1505. log.Printf("[INFO] Setting baseimagename to %#v", baseimagename)
  1506. }
  1507. // Old sane default overrides:
  1508. if baseimageregistry == "ghcr.io" && baseimagename == "shuffle" {
  1509. baseimageregistry = "docker.io"
  1510. baseimagename = "frikky/shuffle"
  1511. os.Setenv("REGISTRY_URL", baseimageregistry)
  1512. os.Setenv("SHUFFLE_BASE_IMAGE_REGISTRY", baseimageregistry)
  1513. os.Setenv("SHUFFLE_BASE_IMAGE_NAME", baseimagename)
  1514. log.Printf("[WARNING] Overriding bad defaults of ghcr.io/shuffle")
  1515. }
  1516. log.Printf("[DEBUG] Setting swarm config to %#v. Default is empty.", swarmConfig)
  1517. // This is now always static
  1518. newWorker := fmt.Sprintf("ghcr.io/shuffle/shuffle-worker:%s", workerVersion)
  1519. if len(newWorkerImage) > 0 {
  1520. newWorker = newWorkerImage
  1521. }
  1522. // Check whether they are the same first
  1523. if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" {
  1524. log.Printf("[DEBUG] Skipping image download as SHUFFLE_AUTO_IMAGE_DOWNLOAD is set to false")
  1525. } else {
  1526. images := []string{
  1527. fmt.Sprintf("frikky/shuffle:app_sdk"),
  1528. newWorker,
  1529. }
  1530. pullOptions := image.PullOptions{}
  1531. for _, image := range images {
  1532. if isKubernetes == "true" {
  1533. log.Printf("[DEBUG] Skipping image pull of '%s' because Kubernetes does it in realtime instead", image)
  1534. } else {
  1535. log.Printf("[DEBUG] Pulling image %s", image)
  1536. reader, err := dockercli.ImagePull(ctx, image, pullOptions)
  1537. if err != nil {
  1538. log.Printf("[ERROR] Failed getting image %s: %s", image, err)
  1539. continue
  1540. }
  1541. io.Copy(os.Stdout, reader)
  1542. log.Printf("[DEBUG] Successfully downloaded and built %s", image)
  1543. }
  1544. }
  1545. }
  1546. }
  1547. func findActiveSwarmNodes() (int64, error) {
  1548. ctx := context.Background()
  1549. nodes, err := dockercli.NodeList(ctx, types.NodeListOptions{})
  1550. if err != nil {
  1551. return 1, err
  1552. }
  1553. nodeCount := int64(0)
  1554. for _, node := range nodes {
  1555. //log.Printf("ID: %s - %#v", node.ID, node.Status.State)
  1556. if node.Status.State == "ready" {
  1557. nodeCount += 1
  1558. }
  1559. }
  1560. // Check for SHUFFLE_MAX_NODES
  1561. // Make it into a number and check if it's lower than nodeCount
  1562. maxNodesString := os.Getenv("SHUFFLE_MAX_SWARM_NODES")
  1563. if len(maxNodesString) > 0 {
  1564. maxNodes, err := strconv.ParseInt(maxNodesString, 10, 64)
  1565. if err != nil {
  1566. return nodeCount, err
  1567. }
  1568. if nodeCount > maxNodes {
  1569. nodeCount = maxNodes
  1570. }
  1571. }
  1572. return nodeCount, nil
  1573. }
  1574. // Get IP
  1575. func getLocalIP() string {
  1576. addrs, err := net.InterfaceAddrs()
  1577. if err != nil {
  1578. return ""
  1579. }
  1580. for _, address := range addrs {
  1581. // check the address type and if it is not a loopback the display it
  1582. if ipnet, ok := address.(*net.IPNet); ok && !ipnet.IP.IsLoopback() {
  1583. if ipnet.IP.To4() != nil {
  1584. return ipnet.IP.String()
  1585. }
  1586. }
  1587. }
  1588. return ""
  1589. }
  1590. // Get all local IPs in the system
  1591. func getLocalIPs() ([]string, error) {
  1592. var ipv4s []string
  1593. var ipv6s []string
  1594. ifaces, err := net.Interfaces()
  1595. if err != nil {
  1596. return nil, err
  1597. }
  1598. for _, iface := range ifaces {
  1599. if iface.Flags&net.FlagUp == 0 {
  1600. continue
  1601. }
  1602. if iface.Flags&net.FlagLoopback != 0 {
  1603. continue
  1604. }
  1605. addrs, err := iface.Addrs()
  1606. if err != nil {
  1607. continue
  1608. }
  1609. for _, address := range addrs {
  1610. ipnet, ok := address.(*net.IPNet)
  1611. if !ok || ipnet.IP == nil {
  1612. continue
  1613. }
  1614. ip := ipnet.IP
  1615. if ip.IsLoopback() || ip.IsLinkLocalUnicast() || ip.IsLinkLocalMulticast() {
  1616. continue
  1617. }
  1618. if ip4 := ip.To4(); ip4 != nil {
  1619. ipv4s = append(ipv4s, ip4.String())
  1620. continue
  1621. }
  1622. if ip.To16() != nil {
  1623. ipv6s = append(ipv6s, ip.String())
  1624. }
  1625. }
  1626. }
  1627. return append(ipv4s, ipv6s...), nil
  1628. }
  1629. func checkSwarmService(ctx context.Context) {
  1630. // https://docs.docker.com/engine/reference/commandline/swarm_init/
  1631. ip := getLocalIP()
  1632. log.Printf("[DEBUG] Attempting swarm setup on %s", ip)
  1633. info, err := dockercli.Info(ctx)
  1634. if err != nil {
  1635. log.Printf("[WARNING] Failed to get Docker Info: %s", err)
  1636. }
  1637. if info.Swarm.ControlAvailable {
  1638. log.Printf("[INFO] Already part of swarm as a manager")
  1639. return
  1640. }
  1641. listenAddr := "0.0.0.0"
  1642. req := swarm.InitRequest{
  1643. ListenAddr: fmt.Sprintf("%s:2377", listenAddr),
  1644. AdvertiseAddr: fmt.Sprintf("%s:2377", ip),
  1645. }
  1646. id, err := dockercli.SwarmInit(ctx, req)
  1647. if err != nil {
  1648. log.Printf("[ERROR] Swarm init issue: %s. Retrying with a failover IP address from interface.", err)
  1649. // Dummy message used for testing
  1650. //err = errors.New("Error response from daemon: could not choose an IP address to advertise since this system has multiple addresses on different interfaces (10.52.208.221 on eno1 and 192.168.122.1 on virbr0) - specify one with --advertise-addr")
  1651. // Update 28 Jan 2026: The error message updated and not as clear
  1652. candidates, err := getLocalIPs()
  1653. if len(candidates) > 0 && err == nil {
  1654. for cnt, candidate := range candidates {
  1655. if cnt > 5 {
  1656. break
  1657. }
  1658. req.AdvertiseAddr = fmt.Sprintf("%s:2377", candidate)
  1659. id, err = dockercli.SwarmInit(context.Background(), req)
  1660. if err != nil {
  1661. continue
  1662. }
  1663. log.Printf("[INFO] Swarm init ID: '%s'.", id)
  1664. return
  1665. }
  1666. }
  1667. log.Printf("[ERROR] Swarm init failed after advertise-addr retries: %s, try running swarm init manually: docker swarm init", err)
  1668. return
  1669. }
  1670. }
  1671. func getContainerResourceUsage(ctx context.Context, cli *dockerclient.Client, containerID string) (float64, float64, error) {
  1672. // Get container stats
  1673. stats, err := cli.ContainerStats(ctx, containerID, false)
  1674. if err != nil {
  1675. return 0, 0, err
  1676. }
  1677. defer stats.Body.Close()
  1678. // Parse and return CPU and memory utilization
  1679. cpuUsage, memoryUsage, err := parseResourceUsage(stats.Body)
  1680. if err != nil {
  1681. return 0, 0, err
  1682. }
  1683. return cpuUsage, memoryUsage, nil
  1684. }
  1685. func parseResourceUsage(body io.Reader) (float64, float64, error) {
  1686. //var stats types.StatsJSON
  1687. var stats container.Stats
  1688. // Decode the stream of stats as JSON
  1689. decoder := json.NewDecoder(body)
  1690. if err := decoder.Decode(&stats); err != nil {
  1691. return 0, 0, err
  1692. }
  1693. //log.Printf("[DEBUG] CPU : %d", stats.CPUStats.CPUUsage.TotalUsage)
  1694. //log.Printf("[DEBUG] CPU2: %d", stats.PreCPUStats.CPUUsage.TotalUsage)
  1695. if stats.CPUStats.CPUUsage.TotalUsage == 0 || stats.PreCPUStats.CPUUsage.TotalUsage == 0 {
  1696. //log.Printf("[DEBUG] BODY: %#v", stats)
  1697. return 0, 0, nil
  1698. }
  1699. // Calculate time difference between current and previous stats in nanoseconds
  1700. timeDelta := float64(stats.Read.Sub(stats.PreRead).Nanoseconds())
  1701. // Calculate CPU usage percentage
  1702. cpuDelta := float64(stats.CPUStats.CPUUsage.TotalUsage - stats.PreCPUStats.CPUUsage.TotalUsage)
  1703. cpuUsage := (cpuDelta / timeDelta) * 100.0
  1704. // Calculate memory usage percentage
  1705. memoryUsage := float64(stats.MemoryStats.Usage) / float64(stats.MemoryStats.Limit) * 100.0
  1706. return cpuUsage, memoryUsage, nil
  1707. }
  1708. func getOrborusStats(ctx context.Context) shuffle.OrborusStats {
  1709. newStats := shuffle.OrborusStats{
  1710. OrgId: org,
  1711. Environment: environment,
  1712. OrborusLabel: orborusLabel,
  1713. Timestamp: time.Now().Unix(),
  1714. Uuid: orborusUuid,
  1715. }
  1716. if (swarmConfig == "run" || swarmConfig == "swarm") && strings.Contains(newWorkerImage, "scale") {
  1717. newStats.Swarm = true
  1718. }
  1719. newStats.PollTime = sleepTime
  1720. newStats.MaxQueue = maxConcurrency
  1721. newStats.Queue = executionCount
  1722. if isKubernetes == "true" || runningMode == "kubernetes" || runningMode == "k8s" {
  1723. newStats.Kubernetes = true
  1724. return newStats
  1725. }
  1726. // Disable orborus stats
  1727. if os.Getenv("SHUFFLE_STATS_DISABLED") == "true" {
  1728. return newStats
  1729. }
  1730. // FIXME: Returning for now due to this causing network congestion
  1731. // and database fillup. The backend api also has it disabled.
  1732. return newStats
  1733. // Use the docker API to get the CPU usage of the docker engine machine
  1734. pers, err := dockercli.Info(ctx)
  1735. if err != nil {
  1736. log.Printf("[ERROR] Failed getting docker info: %s. This is normal IF there are many containers running.", err)
  1737. return newStats
  1738. } else {
  1739. newStats.TotalContainers = pers.Containers
  1740. newStats.StoppedContainers = pers.ContainersStopped
  1741. // Calculate the amount of CPU utilization on the host
  1742. newStats.CPU = int(pers.NCPU)
  1743. newStats.MaxCPU = int(pers.NCPU)
  1744. newStats.Memory = int(pers.MemTotal)
  1745. newStats.MaxMemory = int(pers.MemTotal)
  1746. }
  1747. // Get list of all running containers
  1748. containers, err := dockercli.ContainerList(ctx, container.ListOptions{})
  1749. if err != nil {
  1750. log.Printf("[ERROR] Failed getting container list: %s", err)
  1751. return newStats
  1752. }
  1753. // Use a WaitGroup to wait for all goroutines to finish
  1754. var wg sync.WaitGroup
  1755. // Channel to collect results
  1756. resultCh := make(chan struct {
  1757. containerID string
  1758. cpuUsage float64
  1759. memoryUsage float64
  1760. })
  1761. // Iterate through containers and start a goroutine for each container
  1762. for _, container := range containers {
  1763. // Check if container is running
  1764. if container.State != "running" {
  1765. continue
  1766. }
  1767. wg.Add(1)
  1768. go func(container types.Container) {
  1769. defer wg.Done()
  1770. // Get CPU and memory usage for the container
  1771. cpuUsage, memoryUsage, err := getContainerResourceUsage(ctx, dockercli, container.ID)
  1772. if err != nil {
  1773. //log.Printf("[DEBUG] Error getting resource usage for container %s: %v\n", container.ID, err)
  1774. }
  1775. // Send the result to the channel
  1776. resultCh <- struct {
  1777. containerID string
  1778. cpuUsage float64
  1779. memoryUsage float64
  1780. }{container.ID, cpuUsage, memoryUsage}
  1781. }(container)
  1782. }
  1783. // Close the result channel after all goroutines are done
  1784. go func() {
  1785. wg.Wait()
  1786. close(resultCh)
  1787. }()
  1788. // Collect results from the channel
  1789. // Iterate through containers and get CPU usage
  1790. totalCPU := float64(0.0)
  1791. memUsage := float64(0.0)
  1792. for result := range resultCh {
  1793. //log.Printf("[DEBUG] Container %s CPU utilization: %.2f%%, Memory utilization: %.2f%%\n", result.containerID, result.cpuUsage, result.memoryUsage)
  1794. // check if it's NaN or Inf
  1795. if !math.IsNaN(result.cpuUsage) {
  1796. totalCPU += float64(result.cpuUsage)
  1797. }
  1798. if !math.IsNaN(result.memoryUsage) {
  1799. memUsage += float64(result.memoryUsage)
  1800. }
  1801. }
  1802. newStats.CPUPercent = totalCPU / float64(newStats.CPU)
  1803. newStats.MemoryPercent = memUsage
  1804. //log.Printf("[DEBUG] CPU: %.2f, Memory: %.2f", newStats.CPUPercent, newStats.MemoryPercent)
  1805. /*
  1806. cpuPercent, err := cpu.Percent(250*time.Millisecond, false)
  1807. if err == nil && len(cpuPercent) > 0 {
  1808. newStats.CPUPercent = cpuPercent[0]
  1809. }
  1810. //Percent(interval time.Duration, percpu bool) ([]float64, error)
  1811. // Get memory usage
  1812. memory, err := memory.Get()
  1813. if err != nil {
  1814. log.Printf("[ERROR] Failed getting memory stats: %s", err)
  1815. } else {
  1816. newStats.Memory = int(memory.Used)
  1817. newStats.MaxMemory = int(memory.Total)
  1818. }
  1819. */
  1820. // Get disk usage
  1821. /*
  1822. disk, err := disk.Get()
  1823. if err != nil {
  1824. log.Printf("[ERROR] Failed getting disk stats: %s", err)
  1825. } else {
  1826. newStats.Disk = int(disk.Used)
  1827. newStats.MaxDisk = int(disk.Total)
  1828. }
  1829. */
  1830. /*
  1831. // General
  1832. Disk int `json:"disk"`
  1833. // Docker
  1834. AppContainers int `json:"app_containers"`
  1835. WorkerContainers int `json:"worker_containers"`
  1836. TotalContainers int `json:"total_containers"`
  1837. }
  1838. */
  1839. return newStats
  1840. }
  1841. func sendRemoveRequest(client *http.Client, toBeRemoved shuffle.ExecutionRequestWrapper, baseUrl, environment, auth, org string, sleepTime int) error {
  1842. confirmUrl := fmt.Sprintf("%s/api/v1/workflows/queue/confirm", baseUrl)
  1843. data, err := json.Marshal(toBeRemoved)
  1844. if err != nil {
  1845. log.Printf("[WARNING] Failed removal marshalling: %s", err)
  1846. time.Sleep(time.Duration(sleepTime) * time.Second)
  1847. return err
  1848. }
  1849. result, err := http.NewRequest(
  1850. "POST",
  1851. confirmUrl,
  1852. bytes.NewBuffer([]byte(data)),
  1853. )
  1854. if err != nil {
  1855. log.Printf("[ERROR] Failed building confirm request: %s", err)
  1856. time.Sleep(time.Duration(sleepTime) * time.Second)
  1857. return err
  1858. }
  1859. result.Header.Add("Content-Type", "application/json")
  1860. result.Header.Add("Org-Id", environment)
  1861. if len(auth) > 0 {
  1862. result.Header.Add("Authorization", auth)
  1863. }
  1864. if len(org) > 0 {
  1865. result.Header.Add("Org", org)
  1866. }
  1867. if len(orborusLabel) > 0 {
  1868. result.Header.Add("X-Orborus-Label", orborusLabel)
  1869. }
  1870. resultResp, err := client.Do(result)
  1871. if err != nil {
  1872. if !strings.Contains(fmt.Sprintf("%s", err), "timeout") {
  1873. log.Printf("[ERROR] Failed making confirm request: %s", err)
  1874. }
  1875. time.Sleep(time.Duration(sleepTime) * time.Second)
  1876. return err
  1877. }
  1878. defer resultResp.Body.Close()
  1879. body, err := ioutil.ReadAll(resultResp.Body)
  1880. if err != nil {
  1881. log.Printf("[ERROR] Failed reading confirm body: %s", err)
  1882. time.Sleep(time.Duration(sleepTime) * time.Second)
  1883. return err
  1884. }
  1885. _ = body
  1886. //log.Printf("[DEBUG] Confirm response: %s", string(body))
  1887. return nil
  1888. }
  1889. func cleanup() {
  1890. log.Printf("[INFO] Cleaning up during shutdown")
  1891. ctx := context.Background()
  1892. cleanupExistingNodes(ctx)
  1893. zombiecheck(ctx, 600)
  1894. os.Exit(0)
  1895. }
  1896. func StartAgent() {
  1897. log.Printf("[INFO] Starting Orborus agent mode")
  1898. auditLogEnabled := os.Getenv("SHUFFLE_AUDIT_LOG_ENABLED") == "true"
  1899. if auditLogEnabled {
  1900. log.Printf("[INFO] Audit log monitoring is enabled")
  1901. // Initialize telemetry configuration
  1902. telemetryConfig := shuffle.TelemetryConfig{
  1903. Enabled: true,
  1904. Modes: []string{"audit_log"},
  1905. BufferSize: 1000,
  1906. FlushInterval: 10 * time.Second,
  1907. }
  1908. if excludePatterns := os.Getenv("SHUFFLE_AUDIT_LOG_EXCLUDE"); excludePatterns != "" {
  1909. patterns := strings.Split(excludePatterns, ",")
  1910. telemetryConfig.Filters = append(telemetryConfig.Filters, shuffle.TelemetryFilter{
  1911. Type: "message",
  1912. Exclude: patterns,
  1913. })
  1914. }
  1915. if includePatterns := os.Getenv("SHUFFLE_AUDIT_LOG_INCLUDE"); includePatterns != "" {
  1916. patterns := strings.Split(includePatterns, ",")
  1917. telemetryConfig.Filters = append(telemetryConfig.Filters, shuffle.TelemetryFilter{
  1918. Type: "message",
  1919. Include: patterns,
  1920. })
  1921. }
  1922. collector, err := shuffle.NewAuditLogCollector(telemetryConfig)
  1923. if err != nil {
  1924. log.Printf("[ERROR] Failed to create audit log collector: %v", err)
  1925. } else {
  1926. ctx := context.Background()
  1927. if err := collector.LogCollectorStart(ctx); err != nil {
  1928. log.Printf("[ERROR] Failed to start audit log collector: %v", err)
  1929. } else {
  1930. log.Printf("[INFO] Audit log collector started successfully")
  1931. sigChan := make(chan os.Signal, 1)
  1932. signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
  1933. go func() {
  1934. <-sigChan
  1935. log.Printf("[INFO] Received shutdown signal, stopping audit log collector...")
  1936. collector.Stop()
  1937. os.Exit(0)
  1938. }()
  1939. }
  1940. }
  1941. } else {
  1942. log.Printf("[INFO] Audit log monitoring is disabled")
  1943. }
  1944. select {}
  1945. }
  1946. // Initial loop etc
  1947. func main() {
  1948. // Get arch. amd64 or arm64
  1949. //sigCh := make(chan os.Signal, 1)
  1950. //signal.Notify(sigCh, os.Interrupt, syscall.SIGTERM)
  1951. //defer cleanup()
  1952. agentMode := os.Getenv("SHUFFLE_AGENT_MODE")
  1953. if agentMode == "true" {
  1954. log.Printf("[INFO] Running in agent mode. Starting the agent.")
  1955. StartAgent()
  1956. return
  1957. }
  1958. if os.Getenv("SHUFFLE_PIPELINE_STANDALONE") == "true" {
  1959. log.Printf("[INFO] Allowing use of standalone pipeline (tenzir). URL: %s", pipelineUrl)
  1960. tenzirDisabled = false
  1961. os.Setenv("SHUFFLE_SKIP_PIPELINES", "false")
  1962. os.Setenv("SHUFFLE_PIPELINE_ENABLED", "true")
  1963. }
  1964. // Block until a signal is received
  1965. if shuffle.IsRunningInCluster() {
  1966. log.Printf("[INFO] Running inside k8s cluster")
  1967. }
  1968. if isKubernetes == "true" {
  1969. fixk8sRoles()
  1970. }
  1971. startupDelay := os.Getenv("SHUFFLE_ORBORUS_STARTUP_DELAY")
  1972. if len(startupDelay) > 0 {
  1973. log.Printf("[DEBUG] Setting startup delay to %#v", startupDelay)
  1974. tmpInt, err := strconv.Atoi(startupDelay)
  1975. if err == nil {
  1976. time.Sleep(time.Duration(tmpInt) * time.Second)
  1977. } else {
  1978. log.Printf("[WARNING] Env SHUFFLE_ORBORUS_STARTUP_DELAY must be a number, not '%s'. Using default.", startupDelay)
  1979. }
  1980. }
  1981. // Auto enables pipelines IF they are not mentioned
  1982. if len(os.Getenv("SHUFFLE_SKIP_PIPELINES")) == 0 {
  1983. tenzirDisabled = false
  1984. os.Setenv("SHUFFLE_SKIP_PIPELINES", "false")
  1985. os.Setenv("SHUFFLE_PIPELINE_ENABLED", "true")
  1986. }
  1987. if os.Getenv("SHUFFLE_SKIP_PIPELINES") != "true" && os.Getenv("SHUFFLE_PIPELINE_ENABLED") != "false" {
  1988. // Run in 15 seconds in a goroutine
  1989. go func() {
  1990. time.Sleep(15 * time.Second)
  1991. log.Printf("[INFO] Auto-downloading Sigma rules during startup")
  1992. ruleType := "sigma"
  1993. err := handleFileCategoryChange(ruleType)
  1994. if err != nil {
  1995. log.Printf("[WARNING] Failed downloading %s rules: %s", ruleType, err)
  1996. }
  1997. }()
  1998. }
  1999. log.Println("[INFO] Setting up execution environment for env '%s'", environment)
  2000. // //FIXME
  2001. if baseUrl == "" {
  2002. baseUrl = "https://shuffler.io"
  2003. }
  2004. if len(orborusUuid) == 0 {
  2005. orborusUuid = uuid.NewV4().String()
  2006. }
  2007. //if orgId == "" {
  2008. // log.Printf("[ERROR] Org not defined. Set variable ORG_ID based on your org")
  2009. // os.Exit(3)
  2010. //}
  2011. if environment == "" {
  2012. log.Printf("[ERROR] Environment not defined. Set variable ENVIRONMENT_NAME to configure it.")
  2013. os.Exit(3)
  2014. }
  2015. if timezone == "" {
  2016. timezone = "Europe/Amsterdam"
  2017. }
  2018. log.Printf("[INFO] Using environment '%s' with timezone %s", environment, timezone)
  2019. if len(os.Getenv("SHUFFLE_ORBORUS_PULL_TIME")) > 0 {
  2020. log.Printf("[INFO] Trying to set Orborus sleep time between polls to %s", os.Getenv("SHUFFLE_ORBORUS_PULL_TIME"))
  2021. tmpInt, err := strconv.Atoi(os.Getenv("SHUFFLE_ORBORUS_PULL_TIME"))
  2022. if err == nil {
  2023. sleepTime = tmpInt
  2024. }
  2025. }
  2026. // Handle Cleanup - made it cleanup by default
  2027. if strings.ToLower(os.Getenv("SHUFFLE_CONTAINER_AUTO_CLEANUP")) != "false" && os.Getenv("CLEANUP") == "" {
  2028. cleanupEnv = "true"
  2029. }
  2030. if len(cleanupEnv) > 0 {
  2031. log.Printf("[DEBUG] Verbose mode. NOT cleaning up. Cleanup env: %s", cleanupEnv)
  2032. }
  2033. // Default to 120 instead of default 30
  2034. if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) == 0 {
  2035. os.Setenv("SHUFFLE_APP_SDK_TIMEOUT", "120")
  2036. }
  2037. workerTimeout := 600
  2038. if workerTimeoutEnv != "" {
  2039. tmpInt, err := strconv.Atoi(workerTimeoutEnv)
  2040. if err == nil {
  2041. workerTimeout = tmpInt
  2042. } else {
  2043. log.Printf("[WARNING] Env SHUFFLE_ORBORUS_EXECUTION_TIMEOUT must be a number, not %s", workerTimeoutEnv)
  2044. }
  2045. log.Printf("[INFO] Cleanup process running every %d seconds", workerTimeout)
  2046. }
  2047. if concurrencyEnv != "" {
  2048. //var concurrencyEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY")
  2049. tmpInt, err := strconv.Atoi(concurrencyEnv)
  2050. if err == nil {
  2051. maxConcurrency = tmpInt
  2052. log.Printf("[INFO] Max workflow execution concurrency set to %d", maxConcurrency)
  2053. } else {
  2054. log.Printf("[WARNING] Env SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY must be a number, not %s. Defaulted to %d", workerTimeoutEnv, maxConcurrency)
  2055. }
  2056. }
  2057. if len(os.Getenv("DOCKER_HOST")) > 0 {
  2058. log.Printf("[DEBUG] Running docker with socket proxy %s instead of default", os.Getenv("DOCKER_HOST"))
  2059. } else {
  2060. log.Printf(`[DEBUG] Running docker with default socket /var/run/docker.sock or `)
  2061. }
  2062. ctx := context.Background()
  2063. // Run by default from now
  2064. //commenting for now as its stoppoing minikube
  2065. log.Printf("[INFO] Running towards %s (BASE_URL) with environment name %s", baseUrl, environment)
  2066. if environment == "" {
  2067. environment = "onprem"
  2068. log.Printf("[WARNING] Defaulting to environment name %s. Set environment variable ENVIRONMENT_NAME to change. This should be the same as in the frontend action.", environment)
  2069. }
  2070. if pipelineUrl == "" {
  2071. pipelineUrl = "http://localhost:5160"
  2072. // Find the IP in baseUrl. Base format is http://<ip>:<port>
  2073. if baseUrl != "" && !strings.Contains(baseUrl, "shuffle") && !strings.Contains(baseUrl, "localhost") && !strings.Contains(baseUrl, "run.app") {
  2074. urlSplit := strings.Split(baseUrl, "://")
  2075. if len(urlSplit) > 1 {
  2076. // Find the IP
  2077. ipSplit := strings.Split(urlSplit[1], ":")
  2078. if len(ipSplit) > 0 {
  2079. pipelineUrl = fmt.Sprintf("http://%s:5160", ipSplit[0])
  2080. }
  2081. }
  2082. }
  2083. if len(containerId) > 0 {
  2084. pipelineUrl = "http://tenzir-node:5160"
  2085. }
  2086. log.Printf("[WARNING] SHUFFLE_PIPELINE_URL not set, falling back to default URL: %s. If BASE_URL is set, we use the external IP for that", pipelineUrl)
  2087. os.Setenv("SHUFFLE_PIPELINE_URL", pipelineUrl)
  2088. }
  2089. // FIXME - during init, BUILD and/or LOAD worker and app_sdk
  2090. // Build/load app_sdk so it can be loaded as 127.0.0.1:5000/walkoff_app_sdk
  2091. log.Printf("[INFO] Setting up Docker environment. Downloading worker and App SDK!")
  2092. initializeImages()
  2093. workerImage := fmt.Sprintf("ghcr.io/shuffle/shuffle-worker:%s", workerVersion)
  2094. if len(newWorkerImage) > 0 {
  2095. workerImage = newWorkerImage
  2096. }
  2097. if swarmConfig == "run" || swarmConfig == "swarm" || isKubernetes == "true" {
  2098. if isKubernetes != "true" {
  2099. checkSwarmService(ctx)
  2100. }
  2101. log.Printf("[DEBUG] Cleaning up containers from previous run")
  2102. cleanupExistingNodes(ctx)
  2103. time.Sleep(time.Duration(5) * time.Second)
  2104. log.Printf("[DEBUG] Deploying worker image %s to swarm", workerImage)
  2105. runString := "Run: \"docker service ls\" for more info"
  2106. if isKubernetes != "true" {
  2107. deployServiceWorkers(workerImage)
  2108. err := setBackendToSwarmNetwork(ctx)
  2109. if err != nil {
  2110. log.Printf("[WARNING] Failed setting backend to swarm network: %s", err)
  2111. }
  2112. } else {
  2113. deployK8sWorker(workerImage, "shuffle-workers", []string{})
  2114. runString = "Run: \"kubectl get pods\" for more info"
  2115. }
  2116. log.Printf("[DEBUG] Waiting 45 seconds to ensure workers are deployed. %s", runString)
  2117. time.Sleep(time.Duration(45) * time.Second)
  2118. //deployServiceWorkers(workerImage)
  2119. }
  2120. zombiecheck(ctx, workerTimeout)
  2121. client := shuffle.GetExternalClient(baseUrl)
  2122. fullUrl := fmt.Sprintf("%s/api/v1/workflows/queue", baseUrl)
  2123. // Increases default concurrency to 50 for swarm
  2124. if maxConcurrency < 50 && (swarmConfig == "run" || swarmConfig == "swarm") {
  2125. fullUrl += "?amount=50"
  2126. }
  2127. if isKubernetes == "true" {
  2128. log.Printf("[INFO] Finished configuring kubernetes environment. Connecting to %s", fullUrl)
  2129. } else {
  2130. log.Printf("[INFO] Finished configuring docker environment. Connecting to %s", fullUrl)
  2131. }
  2132. forwardData := bytes.NewBuffer([]byte{})
  2133. forwardMethod := "POST"
  2134. req, err := http.NewRequest(
  2135. forwardMethod,
  2136. fullUrl,
  2137. forwardData,
  2138. )
  2139. if err != nil {
  2140. log.Printf("[ERROR] Failed making request builder during init: %s", err)
  2141. return
  2142. }
  2143. zombiecounter := 0
  2144. req.Header.Add("Content-Type", "application/json")
  2145. req.Header.Add("Org-Id", environment)
  2146. if len(auth) > 0 {
  2147. req.Header.Add("Authorization", auth)
  2148. }
  2149. if len(org) > 0 {
  2150. req.Header.Add("Org", org)
  2151. }
  2152. if len(orborusLabel) > 0 {
  2153. log.Printf("[DEBUG] Sending with Label '%s'", orborusLabel)
  2154. req.Header.Add("X-Orborus-Label", orborusLabel)
  2155. }
  2156. if swarmConfig != "run" && swarmConfig != "swarm" {
  2157. req.Header.Add("X-Orborus-Runmode", "Default")
  2158. } else {
  2159. req.Header.Add("X-Orborus-Runmode", "Docker Swarm")
  2160. }
  2161. if os.Getenv("SHUFFLE_MAX_CPU") != "" {
  2162. // parse
  2163. tmpInt, err := strconv.Atoi(os.Getenv("SHUFFLE_MAX_CPU"))
  2164. if err == nil {
  2165. maxCPUPercent = tmpInt
  2166. }
  2167. }
  2168. swarmPollingTime := time.Now()
  2169. swarmRequestsMade := 0
  2170. swarmControlMode := false
  2171. if os.Getenv("SHUFFLE_SWARM_CONTROL_MODE") == "true" {
  2172. swarmControlMode = true
  2173. }
  2174. log.Printf("[INFO] Waiting for executions at %s with Environment %#v", fullUrl, environment)
  2175. hasStarted := false
  2176. for {
  2177. if req.Method == "POST" {
  2178. // Should find data to send (memory etc.)
  2179. // Create timeout of max 4 seconds just in case
  2180. ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
  2181. defer cancel()
  2182. // Marshal and set body
  2183. orborusStats := getOrborusStats(ctx)
  2184. pipelinePayload, pipelineerr := sendPipelineHealthStatus()
  2185. if pipelineerr != nil {
  2186. // Too verbose to be enabled.
  2187. //log.Printf("[ERROR] Failed sending pipeline health status: %s", pipelineerr)
  2188. }
  2189. orborusStats.DataLake = pipelinePayload
  2190. jsonData, err := json.Marshal(orborusStats)
  2191. if err == nil {
  2192. req.Body = ioutil.NopCloser(bytes.NewBuffer(jsonData))
  2193. } else {
  2194. log.Printf("[ERROR] Failed marshalling. Maybe max 4 second timeout? %s", err)
  2195. }
  2196. if int(orborusStats.CPUPercent) > maxCPUPercent {
  2197. log.Printf("[DEBUG] CPU usage is at %f%%. This is more than the max limit the machine should be running at (%d). Waiting before continue.", orborusStats.CPUPercent, maxCPUPercent)
  2198. time.Sleep(time.Duration(sleepTime) * time.Second)
  2199. continue
  2200. }
  2201. }
  2202. newresp, err := client.Do(req)
  2203. if err != nil {
  2204. log.Printf("[WARNING] Failed making request to %s: %s", fullUrl, err)
  2205. zombiecounter += 1
  2206. if zombiecounter*sleepTime > workerTimeout {
  2207. go zombiecheck(ctx, workerTimeout)
  2208. zombiecounter = 0
  2209. }
  2210. time.Sleep(time.Duration(sleepTime) * time.Second)
  2211. continue
  2212. }
  2213. //defer newresp.Body.Close()
  2214. if newresp.StatusCode == 405 {
  2215. log.Printf("[WARNING] Received 405 from %s. This is likely due to a misconfigured base URL. Automatically swapping to GET request (backwards compatibility)", fullUrl)
  2216. req.Method = "GET"
  2217. req.Body = nil
  2218. //time.Sleep(time.Duration(sleepTime) * time.Second)
  2219. continue
  2220. }
  2221. body, err := ioutil.ReadAll(newresp.Body)
  2222. if err != nil {
  2223. log.Printf("[ERROR] Failed reading body from Shuffle: %s", err)
  2224. zombiecounter += 1
  2225. if zombiecounter*sleepTime > workerTimeout {
  2226. go zombiecheck(ctx, workerTimeout)
  2227. zombiecounter = 0
  2228. }
  2229. time.Sleep(time.Duration(sleepTime) * time.Second)
  2230. continue
  2231. }
  2232. // Controls Leader/Follower mode
  2233. if newresp.StatusCode == 409 {
  2234. log.Printf("[INFO] Another Orborus is already handling jobs. Polling every 30 seconds in case Leader stops. Resp: %s", string(body))
  2235. time.Sleep(time.Duration(30) * time.Second)
  2236. continue
  2237. } else if newresp.StatusCode != 200 {
  2238. log.Printf("[ERROR] Backend connection failed for url '%s', or is missing (%d): %s", fullUrl, newresp.StatusCode, string(body))
  2239. } else {
  2240. if !hasStarted {
  2241. log.Printf("[DEBUG] Starting iteration on environment %#v (default = Shuffle). Got statuscode %d from backend on first request", environment, newresp.StatusCode)
  2242. }
  2243. if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" && os.Getenv("SHUFFLE_SCALE_REPLICAS") == "" {
  2244. //go AutoScale(ctx)
  2245. }
  2246. hasStarted = true
  2247. }
  2248. var executionRequests shuffle.ExecutionRequestWrapper
  2249. err = json.Unmarshal(body, &executionRequests)
  2250. if err != nil {
  2251. log.Printf("[WARNING] Failed executionrequest in queue unmarshaling: %s", err)
  2252. sleepTime = 10
  2253. zombiecounter += 1
  2254. if zombiecounter*sleepTime > workerTimeout {
  2255. go zombiecheck(ctx, workerTimeout)
  2256. zombiecounter = 0
  2257. }
  2258. time.Sleep(time.Duration(sleepTime) * time.Second)
  2259. continue
  2260. }
  2261. if hasStarted && len(executionRequests.Data) > 0 {
  2262. //log.Printf("[INFO] Body: %s", string(body))
  2263. // Type string `json:"type"`
  2264. }
  2265. // FIXME: Add features here for orborus & worker to
  2266. // do things on behalf of backend
  2267. var toBeRemoved shuffle.ExecutionRequestWrapper
  2268. if len(executionRequests.Data) > 0 {
  2269. newrequests := []shuffle.ExecutionRequest{}
  2270. // Deduplicating in case same job shows up multiple times
  2271. // This is specifically to handle data pipelines better
  2272. deduplicatedJobs := []shuffle.ExecutionRequest{}
  2273. for _, incRequest := range executionRequests.Data {
  2274. if !strings.Contains(incRequest.Type, "DOCKER") && !strings.Contains(incRequest.Type, "PIPELINE") && !strings.Contains(incRequest.Type, "SIGMA") && !strings.Contains(incRequest.Type, "TENZIR") {
  2275. deduplicatedJobs = append(deduplicatedJobs, incRequest)
  2276. continue
  2277. }
  2278. found := false
  2279. for _, dedupRequest := range deduplicatedJobs {
  2280. if incRequest.ExecutionArgument == dedupRequest.ExecutionArgument && incRequest.Type == dedupRequest.Type {
  2281. found = true
  2282. break
  2283. }
  2284. }
  2285. if found {
  2286. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2287. continue
  2288. }
  2289. deduplicatedJobs = append(deduplicatedJobs, incRequest)
  2290. }
  2291. executionRequests.Data = deduplicatedJobs
  2292. for _, incRequest := range executionRequests.Data {
  2293. // Looking for specific jobs
  2294. if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_START" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" || incRequest.Type == "PIPELINE_UPDATE" {
  2295. log.Printf("[INFO] Handling pipeline request from backend: '%s' with argument '%s'", incRequest.Type, incRequest.ExecutionArgument)
  2296. os.Setenv("SHUFFLE_SKIP_PIPELINES", "false")
  2297. os.Setenv("SHUFFLE_PIPELINE_ENABLED", "true")
  2298. tenzirDisabled = false
  2299. // Running NEW or editing pipelines
  2300. err := handlePipeline(incRequest)
  2301. if err != nil {
  2302. log.Printf("[ERROR] Failed handling pipeline ('%s' '%s'): %s. Deleting job anyway.", incRequest.Type, incRequest.ExecutionSource, err)
  2303. }
  2304. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2305. } else if incRequest.Type == "DOCKER_IMAGE_DOWNLOAD" {
  2306. log.Printf("[INFO] Re-downloading new image(s) due to backend request: %#v", incRequest.ExecutionArgument)
  2307. if len(incRequest.ExecutionArgument) > 0 {
  2308. go handleBackendImageDownload(ctx, incRequest.ExecutionArgument)
  2309. } else {
  2310. log.Printf("[ERROR] No image name provided for download. Removing job from queue.")
  2311. }
  2312. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2313. } else if incRequest.Type == "CATEGORY_UPDATE" {
  2314. os.Setenv("SHUFFLE_SKIP_PIPELINES", "false")
  2315. tenzirDisabled = false
  2316. err = handleFileCategoryChange("sigma")
  2317. if err != nil {
  2318. log.Printf("[ERROR] Failed to download the file category: %s", err)
  2319. }
  2320. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2321. } else if incRequest.Type == "DISABLE_SIGMA_FOLDER" {
  2322. log.Printf("[INFO] Got job to disable sigma rules")
  2323. err = removeFileCategory("sigma")
  2324. if err != nil {
  2325. log.Printf("[ERROR] Failed to disable the sigma rules: %s", err)
  2326. }
  2327. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2328. } else if incRequest.Type == "DISABLE_SIGMA_FILE" {
  2329. fileName := incRequest.ExecutionArgument
  2330. log.Printf("[INFO] Got job to disable sigma file %s", fileName)
  2331. err = disableRule(fileName)
  2332. if err != nil {
  2333. log.Printf("[ERROR] Failed to disable the sigma file %s, reason: %s", fileName, err)
  2334. }
  2335. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2336. } else if incRequest.Type == "ENABLE_SIGMA_FILE" {
  2337. fileName := incRequest.ExecutionArgument
  2338. log.Printf("[INFO] Got job to enable sigma file %s", fileName)
  2339. err = enableRule(fileName)
  2340. if err != nil {
  2341. log.Printf("[ERROR] Failed to disable the sigma file %s, reason: %s", fileName, err)
  2342. }
  2343. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2344. } else if incRequest.Type == "START_TENZIR" {
  2345. log.Printf("[INFO] Got job to start tenzir")
  2346. // Manual command = overrides to allow starting of Tenzir from the frontend anyway.
  2347. //os.Setenv("SHUFFLE_SKIP_PIPELINES", "false")
  2348. tenzirDisabled = false
  2349. // Removed either way
  2350. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2351. err := deployTenzirNode()
  2352. if err != nil {
  2353. if strings.Contains(fmt.Sprintf("%s", err), "node available") {
  2354. // Disabling until UI is updated
  2355. //os.Setenv("SHUFFLE_SKIP_PIPELINES", "true")
  2356. //tenzirDisabled = true
  2357. log.Printf("[ERROR] Failed to start tenzir, reason: %s", err)
  2358. err = shuffle.CreateOrgNotification(
  2359. ctx,
  2360. fmt.Sprintf("Failed to start Tenzir: %s", err),
  2361. fmt.Sprintf("Tenzir failed to start due to: %s", err),
  2362. fmt.Sprintf("/detections/Sigma"),
  2363. org,
  2364. true,
  2365. "LOW",
  2366. "TENZIR_START",
  2367. )
  2368. if err != nil {
  2369. log.Printf("[ERROR] Failed to send notification: %s", err)
  2370. return
  2371. }
  2372. }
  2373. }
  2374. } else {
  2375. if debug {
  2376. log.Printf("[DEBUG] Passing execution ID request to normal queue: %#v", incRequest.ExecutionId)
  2377. }
  2378. newrequests = append(newrequests, incRequest)
  2379. }
  2380. }
  2381. if len(toBeRemoved.Data) > 0 {
  2382. err = sendRemoveRequest(client, toBeRemoved, baseUrl, environment, auth, org, sleepTime)
  2383. if err != nil {
  2384. log.Printf("[ERROR] Failed sending remove request: %s", err)
  2385. } else {
  2386. toBeRemoved.Data = []shuffle.ExecutionRequest{}
  2387. }
  2388. }
  2389. // Remove the download image request
  2390. executionRequests.Data = newrequests
  2391. }
  2392. // Skipping throttling with swarm
  2393. if swarmConfig != "run" && swarmConfig != "swarm" {
  2394. if len(executionRequests.Data) == 0 {
  2395. zombiecounter += 1
  2396. if zombiecounter*sleepTime > workerTimeout {
  2397. go zombiecheck(ctx, workerTimeout)
  2398. zombiecounter = 0
  2399. }
  2400. time.Sleep(time.Duration(sleepTime) * time.Second)
  2401. continue
  2402. }
  2403. // Anything below here verifies concurrency
  2404. executionCount = getRunningWorkers(ctx, workerTimeout)
  2405. if executionCount >= maxConcurrency {
  2406. if zombiecounter*sleepTime > workerTimeout {
  2407. go zombiecheck(ctx, workerTimeout)
  2408. zombiecounter = 0
  2409. }
  2410. time.Sleep(time.Duration(sleepTime) * time.Second)
  2411. continue
  2412. }
  2413. allowed := maxConcurrency - executionCount
  2414. if len(executionRequests.Data) > allowed {
  2415. log.Printf("[WARNING] Throttle - Cutting down requests from %d to %d (MAX: %d, CUR: %d)", len(executionRequests.Data), allowed, maxConcurrency, executionCount)
  2416. executionRequests.Data = executionRequests.Data[0:allowed]
  2417. }
  2418. } else if swarmControlMode && (swarmConfig == "run" || swarmConfig == "swarm") {
  2419. // any reason it is not maxConcurrency instead of
  2420. // hardcoded 50?
  2421. if len(executionRequests.Data) > 50 {
  2422. executionRequests.Data = executionRequests.Data[0:50]
  2423. }
  2424. if swarmRequestsMade > 100 && time.Since(swarmPollingTime).Seconds() > 5 {
  2425. log.Printf("[DEBUG] Swarm requests made: %d", swarmRequestsMade)
  2426. time.Sleep(time.Duration(sleepTime) * time.Second)
  2427. swarmPollingTime = time.Now()
  2428. swarmRequestsMade = 0
  2429. }
  2430. swarmRequestsMade += len(executionRequests.Data)
  2431. }
  2432. // New, abortable version. Should check executionid and remove everything else
  2433. for _, execution := range executionRequests.Data {
  2434. if len(execution.ExecutionArgument) > 0 {
  2435. log.Printf("[INFO] Argument: %s", execution.ExecutionArgument)
  2436. }
  2437. if execution.Type == "schedule" {
  2438. log.Printf("[INFO] Schedule type! Weird deployment. Type: %s", execution.Type)
  2439. continue
  2440. }
  2441. if len(execution.ExecutionId) == 0 {
  2442. log.Printf("[WARNING] Execution ID is empty: %#v", execution)
  2443. continue
  2444. }
  2445. if execution.Status == "ABORT" || execution.Status == "FAILED" {
  2446. log.Printf("[INFO][%s] Executionstatus issue: ", execution.ExecutionId, execution.Status)
  2447. }
  2448. if shuffle.ArrayContains(executionIds, execution.ExecutionId) {
  2449. log.Printf("[INFO][%s] Execution already handled (rerunning old execution)", execution.ExecutionId)
  2450. toBeRemoved.Data = append(toBeRemoved.Data, execution)
  2451. // Should check when last this was ran, and if it's more than 10 minutes ago and it's not finished, we should run it again?
  2452. /*
  2453. if swarmConfig != "run" && swarmConfig != "swarm" {
  2454. continue
  2455. }
  2456. */
  2457. }
  2458. // Now, how do I execute this one?
  2459. containerName := fmt.Sprintf("worker-%s", execution.ExecutionId)
  2460. env := []string{
  2461. fmt.Sprintf("AUTHORIZATION=%s", execution.Authorization),
  2462. fmt.Sprintf("EXECUTIONID=%s", execution.ExecutionId),
  2463. fmt.Sprintf("ENVIRONMENT_NAME=%s", environment),
  2464. fmt.Sprintf("BASE_URL=%s", baseUrl),
  2465. fmt.Sprintf("CLEANUP=%s", cleanupEnv),
  2466. fmt.Sprintf("TZ=%s", timezone),
  2467. fmt.Sprintf("SHUFFLE_PASS_APP_PROXY=%s", os.Getenv("SHUFFLE_PASS_APP_PROXY")),
  2468. fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", os.Getenv("SHUFFLE_SWARM_CONFIG")),
  2469. fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", os.Getenv("SHUFFLE_LOGS_DISABLED")),
  2470. fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", os.Getenv("SHUFFLE_BASE_IMAGE_NAME")),
  2471. fmt.Sprintf("SHUFFLE_ALLOW_PACKAGE_INSTAL=%s", os.Getenv("SHUFFLE_ALLOW_PACKAGE_INSTALL")),
  2472. }
  2473. //log.Printf("Running worker with proxy? %s", os.Getenv("SHUFFLE_PASS_WORKER_PROXY"))
  2474. if strings.ToLower(os.Getenv("SHUFFLE_PASS_WORKER_PROXY")) == "true" {
  2475. env = append(env, fmt.Sprintf("HTTP_PROXY=%s", os.Getenv("HTTP_PROXY")))
  2476. env = append(env, fmt.Sprintf("HTTPS_PROXY=%s", os.Getenv("HTTPS_PROXY")))
  2477. env = append(env, fmt.Sprintf("NO_PROXY=%s", os.Getenv("NO_PROXY")))
  2478. }
  2479. if dockerApiVersion != "" {
  2480. env = append(env, fmt.Sprintf("DOCKER_API_VERSION=%s", dockerApiVersion))
  2481. }
  2482. if len(os.Getenv("DOCKER_HOST")) > 0 {
  2483. env = append(env, fmt.Sprintf("DOCKER_HOST=%s", os.Getenv("DOCKER_HOST")))
  2484. }
  2485. if len(os.Getenv("SHUFFLE_MEMCACHED")) > 0 {
  2486. env = append(env, fmt.Sprintf("SHUFFLE_MEMCACHED=%s", os.Getenv("SHUFFLE_MEMCACHED")))
  2487. }
  2488. if len(os.Getenv("SHUFFLE_CLOUDRUN_URL")) > 0 {
  2489. env = append(env, fmt.Sprintf("SHUFFLE_CLOUDRUN_URL=%s", os.Getenv("SHUFFLE_CLOUDRUN_URL")))
  2490. }
  2491. if len(os.Getenv("SHUFFLE_SKIPSSL_VERIFY")) > 0 {
  2492. env = append(env, fmt.Sprintf("SHUFFLE_SKIPSSL_VERIFY=%s", os.Getenv("SHUFFLE_SKIPSSL_VERIFY")))
  2493. }
  2494. if len(os.Getenv("SHUFFLE_DEBUG_MEMORY")) > 0 {
  2495. env = append(env, fmt.Sprintf("SHUFFLE_DEBUG_MEMORY=%s", os.Getenv("SHUFFLE_DEBUG_MEMORY")))
  2496. }
  2497. // Look for volume binds
  2498. if len(os.Getenv("SHUFFLE_VOLUME_BINDS")) > 0 {
  2499. //log.Printf("[DEBUG] Added volume binds: %s", os.Getenv("SHUFFLE_VOLUME_BINDS"))
  2500. env = append(env, fmt.Sprintf("SHUFFLE_VOLUME_BINDS=%s", os.Getenv("SHUFFLE_VOLUME_BINDS")))
  2501. }
  2502. if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 {
  2503. env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
  2504. }
  2505. // Setting up internal proxy config for Shuffle -> shuffle comms
  2506. overrideHttpProxy := os.Getenv("SHUFFLE_INTERNAL_HTTP_PROXY")
  2507. overrideHttpsProxy := os.Getenv("SHUFFLE_INTERNAL_HTTPS_PROXY")
  2508. if len(overrideHttpProxy) > 0 {
  2509. env = append(env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTP_PROXY=%s", overrideHttpProxy))
  2510. }
  2511. if len(overrideHttpsProxy) > 0 {
  2512. env = append(env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTPS_PROXY=%s", overrideHttpsProxy))
  2513. }
  2514. if len(os.Getenv("SHUFFLE_MAX_SWARM_NODES")) > 0 {
  2515. env = append(env, fmt.Sprintf("SHUFFLE_MAX_SWARM_NODES=%s", os.Getenv("SHUFFLE_MAX_SWARM_NODES")))
  2516. }
  2517. err = deployWorker(workerImage, containerName, env, execution)
  2518. zombiecounter += 1
  2519. if err == nil {
  2520. //log.Printf("[DEBUG] ExecutionID %s was deployed and to be removed from queue.", execution.ExecutionId)
  2521. toBeRemoved.Data = append(toBeRemoved.Data, execution)
  2522. executionIds = append(executionIds, execution.ExecutionId)
  2523. } else {
  2524. log.Printf("[WARNING][%s] Failed to deploy: %s", execution.ExecutionId, err)
  2525. if strings.Contains(err.Error(), "already exists") {
  2526. toBeRemoved.Data = append(toBeRemoved.Data, execution)
  2527. executionIds = append(executionIds, execution.ExecutionId)
  2528. } else if strings.Contains(err.Error(), "No such image") {
  2529. // Download the image
  2530. if isKubernetes == "true" {
  2531. log.Printf("[DEBUG] Skipping image pull of '%s' because Kubernetes does it in realtime instead", workerImage)
  2532. } else {
  2533. log.Printf("[DEBUG] Re-pulling image %s as it doesn't exist, and is necessary for worker to run (autofix)", workerImage)
  2534. pullOptions := image.PullOptions{}
  2535. _, err = dockercli.ImagePull(ctx, workerImage, pullOptions)
  2536. if err != nil {
  2537. log.Printf("[ERROR] Failed to pull image %s: %s", workerImage, err)
  2538. }
  2539. }
  2540. }
  2541. }
  2542. }
  2543. // Removes handled workflows (worker is made)
  2544. //log.Printf("\n\n[INFO] Removing %d executions from queue\n\n", len(toBeRemoved.Data))
  2545. if len(toBeRemoved.Data) > 0 {
  2546. err = sendRemoveRequest(client, toBeRemoved, baseUrl, environment, auth, org, sleepTime)
  2547. if err != nil {
  2548. log.Printf("[ERROR] Failed to remove executions from queue: %s", err)
  2549. }
  2550. }
  2551. time.Sleep(time.Duration(sleepTime) * time.Second)
  2552. }
  2553. }
  2554. // Tenzir command samples
  2555. // docker pull ghcr.io/dominiklohmann/tenzir-arm64:latest
  2556. // docker tag ghcr.io/dominiklohmann/tenzir-arm64:latest tenzir/tenzir:latest
  2557. // Read from Cache and send it to a webhook
  2558. // docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines'
  2559. func handlePipeline(incRequest shuffle.ExecutionRequest) error {
  2560. log.Printf("[INFO] Pipeline: '%s' with source '%s'", incRequest.Type, incRequest.ExecutionSource)
  2561. err := deployTenzirNode()
  2562. if err != nil {
  2563. log.Printf("[ERROR] Failed to deploy the pipeline, reason: %s", err)
  2564. return err
  2565. }
  2566. // no need of execution arguments for STOP and DELETE
  2567. if (incRequest.Type != "PIPELINE_STOP" && incRequest.Type != "PIPELINE_DELETE") && len(incRequest.ExecutionArgument) == 0 {
  2568. log.Printf("[ERROR] No execution argument found for pipeline type %s. Skipping", incRequest.Type)
  2569. return errors.New("no execution argument found for pipeline create. Skipping")
  2570. }
  2571. identifier := strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-"))
  2572. if !strings.HasPrefix(strings.ToLower(incRequest.ExecutionSource), "shuffle") {
  2573. identifier = fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-")))
  2574. }
  2575. command := incRequest.ExecutionArgument
  2576. pipelines = []shuffle.PipelineInfo{}
  2577. if incRequest.Type == "PIPELINE_CREATE" {
  2578. log.Printf("[INFO] Should delete -> recreate new pipeline with id %#v", identifier)
  2579. //err := deployPipeline(image, identifier, command)
  2580. _, err := createPipeline(command, identifier)
  2581. if err != nil {
  2582. log.Printf("[ERROR] Failed to create pipeline: %s", err)
  2583. return err
  2584. }
  2585. } else if incRequest.Type == "PIPELINE_DELETE" || incRequest.Type == "PIPELINE_STOP" {
  2586. pipelineId := incRequest.ExecutionId
  2587. log.Printf("[INFO] Should delete pipeline %#v. PipelineID: %s", identifier, pipelineId)
  2588. //pipelineId, err := searchPipeline(identifier)
  2589. //if err != nil {
  2590. //}
  2591. err = deletePipeline(pipelineId)
  2592. if err != nil {
  2593. log.Printf("[ERROR] Failed Deleting Pipeline %s", err)
  2594. return err
  2595. }
  2596. /*
  2597. } else if incRequest.Type == "PIPELINE_STOP" {
  2598. log.Printf("[INFO] Should stop the pipeline %#v", identifier)
  2599. pipelineId, err := searchPipeline(identifier)
  2600. if err != nil {
  2601. log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
  2602. toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
  2603. return err
  2604. }
  2605. _, err = updatePipelineState(command, pipelineId, "stop")
  2606. if err != nil {
  2607. log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err)
  2608. return err
  2609. } else {
  2610. log.Printf("[INFO] Successfully stopped the Pipeline: %s", pipelineId)
  2611. }
  2612. */
  2613. } else if incRequest.Type == "PIPELINE_START" {
  2614. log.Printf("[INFO] Should start the pipeline %#v", identifier)
  2615. pipelineId, err := searchPipeline(identifier)
  2616. if err != nil {
  2617. if err.Error() == "no existing pipeline found with name" {
  2618. log.Printf("[INFO] Starting a new pipeline with command '%s' and identifier '%s'", command, identifier)
  2619. var createErr error
  2620. pipelineId, createErr = createPipeline(command, identifier)
  2621. if createErr != nil {
  2622. return createErr
  2623. }
  2624. } else {
  2625. log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
  2626. return err
  2627. }
  2628. }
  2629. log.Printf("[INFO] Starting existing pipeline with ID %s", pipelineId)
  2630. _, err = updatePipelineState(command, pipelineId, "start")
  2631. if err != nil {
  2632. log.Printf("[ERROR] Failed to start Pipeline: %s reason:%s ", pipelineId, err)
  2633. return err
  2634. } else {
  2635. log.Printf("[INFO] Successfully started pipeline: %s", pipelineId)
  2636. }
  2637. } else {
  2638. log.Printf("[ERROR] Unknown type for pipeline: %s", incRequest.Type)
  2639. return errors.New("unknown type for pipeline")
  2640. }
  2641. return nil
  2642. }
  2643. func deployTenzirNode() error {
  2644. // Specifically for standalone tenzir
  2645. if os.Getenv("SHUFFLE_PIPELINE_STANDALONE") == "true" {
  2646. return nil
  2647. }
  2648. // Disabled all pipeline features
  2649. if os.Getenv("SHUFFLE_SKIP_PIPELINES") == "true" {
  2650. return errors.New("Pipelines are disabled by user with SHUFFLE_SKIP_PIPELINES (1)")
  2651. }
  2652. if isKubernetes == "true" {
  2653. return errors.New("Tenzir not implemented for k8s")
  2654. }
  2655. err := checkTenzirNode()
  2656. if err == nil {
  2657. return nil
  2658. }
  2659. ctx := context.Background()
  2660. cacheKey := "tenzir-key"
  2661. _, err = shuffle.GetCache(ctx, cacheKey)
  2662. if err == nil {
  2663. return nil
  2664. }
  2665. imageName := "frikky/shuffle:tenzir"
  2666. if os.Getenv("TENZIR_IMAGE_NAME") != "" {
  2667. imageName = os.Getenv("TENZIR_IMAGE_NAME")
  2668. log.Printf("[INFO] Using custom Tenzir image name: %s", imageName)
  2669. }
  2670. containerName := "tenzir-node"
  2671. containerStartOptions := container.StartOptions{}
  2672. containerInfo, err := dockercli.ContainerInspect(ctx, containerName)
  2673. if err != nil {
  2674. if dockerclient.IsErrNotFound(err) {
  2675. // Create network if it doesn't exist
  2676. networkName := "tenzir-network"
  2677. networkSubnet := "192.168.102.0/24"
  2678. networkGateway := "192.168.102.1"
  2679. err = createNetworkIfNotExists(ctx, networkName, networkSubnet, networkGateway)
  2680. if err != nil {
  2681. log.Printf("[ERROR] Failed to create network %s: %s", networkName, err)
  2682. //return err
  2683. }
  2684. // Trying to connect orborus to the tenzir network as well
  2685. err = dockercli.NetworkConnect(ctx, networkName, containerId, nil)
  2686. if err != nil {
  2687. log.Printf("[ERROR] Error connecting tenzir container to network: %s", err)
  2688. }
  2689. // Check if image exists
  2690. _, _, err := dockercli.ImageInspectWithRaw(ctx, imageName)
  2691. if dockerclient.IsErrNotFound(err) {
  2692. log.Printf("[DEBUG] Pulling image %s. This may take a while.", imageName)
  2693. pullOptions := image.PullOptions{}
  2694. out, err := dockercli.ImagePull(ctx, imageName, pullOptions)
  2695. if err != nil {
  2696. log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err)
  2697. return err
  2698. }
  2699. defer out.Close()
  2700. io.Copy(io.Discard, out)
  2701. } else if err != nil {
  2702. return err
  2703. }
  2704. err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions)
  2705. if err != nil {
  2706. return err
  2707. }
  2708. } else {
  2709. return err
  2710. }
  2711. } else {
  2712. if !containerInfo.State.Running {
  2713. log.Printf("[DEBUG] Tenzir Node exists, but is not running. Restarting it.")
  2714. err := dockercli.ContainerStart(ctx, containerName, containerStartOptions)
  2715. if err != nil {
  2716. log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err)
  2717. return err
  2718. }
  2719. time.Sleep(10 * time.Second)
  2720. log.Printf("[INFO] Waiting for Tenzir to become available ...")
  2721. err = checkTenzirNode()
  2722. if err != nil {
  2723. return err
  2724. }
  2725. }
  2726. }
  2727. tenzirStatus := struct {
  2728. ContainerStatus string `json:"container_status"`
  2729. }{
  2730. ContainerStatus: "running",
  2731. }
  2732. cacheData, err := json.Marshal(tenzirStatus)
  2733. if err != nil {
  2734. log.Printf("[WARNING] Failed marshalling execution: %s", err)
  2735. }
  2736. err = shuffle.SetCache(ctx, cacheKey, cacheData, 1)
  2737. if err != nil {
  2738. log.Printf("[WARNING] Failed updating cache for tenzir: %s", err)
  2739. }
  2740. return nil
  2741. }
  2742. func createAndStartTenzirNode(ctx context.Context, containerName, imageName string, containerStartOptions container.StartOptions) error {
  2743. healthconfig := &container.HealthConfig{
  2744. Test: []string{"tenzir --connection-timeout=30s --connection-retry-delay=1s 'api /ping'"},
  2745. Interval: 30 * time.Second,
  2746. Retries: 1,
  2747. }
  2748. // Ensure restart policy is there
  2749. config := &container.Config{
  2750. Hostname: containerName,
  2751. Cmd: []string{"--commands=web server --mode=dev --bind=0.0.0.0"},
  2752. Image: imageName,
  2753. Healthcheck: healthconfig,
  2754. ExposedPorts: nat.PortSet{
  2755. "5160/tcp": struct{}{},
  2756. "1514/udp": struct{}{},
  2757. "1514/tcp": struct{}{},
  2758. },
  2759. Entrypoint: []string{containerName},
  2760. Env: []string{},
  2761. }
  2762. tenzirApikey := os.Getenv("TENZIR_PLUGINS__PLATFORM__API_KEY")
  2763. tenzirControlEndpoint := os.Getenv("TENZIR_PLUGINS__PLATFORM__CONTROL_ENDPOINT")
  2764. tenzirPluginsPlatform := os.Getenv("TENZIR_PLUGINS__PLATFORM__TENANT_ID")
  2765. anyFound := false
  2766. if len(tenzirApikey) > 0 {
  2767. config.Env = append(config.Env, fmt.Sprintf("TENZIR_PLUGINS__PLATFORM__API_KEY=%s", tenzirApikey))
  2768. anyFound = true
  2769. }
  2770. if len(tenzirControlEndpoint) > 0 {
  2771. config.Env = append(config.Env, fmt.Sprintf("TENZIR_PLUGINS__PLATFORM__CONTROL_ENDPOINT=%s", tenzirControlEndpoint))
  2772. anyFound = true
  2773. }
  2774. if len(tenzirPluginsPlatform) > 0 {
  2775. config.Env = append(config.Env, fmt.Sprintf("TENZIR_PLUGINS__PLATFORM__TENANT_ID=%s", tenzirPluginsPlatform))
  2776. anyFound = true
  2777. }
  2778. tenzirStorageFolder := os.Getenv("SHUFFLE_STORAGE_FOLDER")
  2779. if len(tenzirStorageFolder) > 0 {
  2780. tenzirStorageFolder = tenzirStorageFolder
  2781. if !strings.HasSuffix(tenzirStorageFolder, "/") {
  2782. tenzirStorageFolder = tenzirStorageFolder + "/"
  2783. }
  2784. } else {
  2785. tenzirStorageFolder = "/tmp/"
  2786. log.Printf("[DEBUG] Using base folder %s for Tenzir storage. Change it using environment variable SHUFFLE_STORAGE_FOLDER=/filepath/", tenzirStorageFolder)
  2787. }
  2788. if !anyFound {
  2789. //log.Printf("[DEBUG] No Tenzir Plugin environment variables found.")
  2790. } else {
  2791. //log.Printf("[DEBUG] Attempting Tenzir connection with app.tenzir.com tenant '%s'", tenzirPluginsPlatform)
  2792. }
  2793. hostConfig := &container.HostConfig{
  2794. PortBindings: nat.PortMap{
  2795. "1514/tcp": []nat.PortBinding{{HostPort: "1514"}},
  2796. "1514/udp": []nat.PortBinding{{HostPort: "1514"}},
  2797. "5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
  2798. },
  2799. Mounts: []mount.Mount{
  2800. {
  2801. Type: "bind",
  2802. Source: tenzirStorageFolder,
  2803. Target: "/tmp",
  2804. },
  2805. /*
  2806. {
  2807. Type: "bind",
  2808. Source: tenzirStorageFolder,
  2809. Target: "/var/log/tenzir/",
  2810. },
  2811. {
  2812. Type: "bind",
  2813. Source: tenzirStorageFolder,
  2814. Target: "/var/cache/tenzir/",
  2815. },
  2816. */
  2817. },
  2818. VolumeDriver: "local",
  2819. RestartPolicy: container.RestartPolicy{
  2820. Name: "always",
  2821. },
  2822. }
  2823. if os.Getenv("SHUFFLE_DISABLE_SYSLOG") == "true" {
  2824. hostConfig.PortBindings = nat.PortMap{
  2825. "5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
  2826. }
  2827. }
  2828. if skipPipelineMount {
  2829. hostConfig.Mounts = []mount.Mount{}
  2830. }
  2831. //networkingConfig := &network.NetworkingConfig{
  2832. // EndpointsConfig: map[string]*network.EndpointSettings{
  2833. // "tenzir-network": {
  2834. // IPAMConfig: &network.EndpointIPAMConfig{
  2835. // IPv4Address: "192.168.102.100",
  2836. // },
  2837. // },
  2838. // },
  2839. //}
  2840. networkingConfig := &network.NetworkingConfig{
  2841. EndpointsConfig: map[string]*network.EndpointSettings{
  2842. "tenzir-network": {
  2843. IPAMConfig: nil,
  2844. Aliases: []string{"tenzir-node"},
  2845. },
  2846. },
  2847. }
  2848. // FIXME: Is this necessary? Seems to screw up networking:
  2849. // conflicting options: hostname and the network mode
  2850. /*
  2851. if isKubernetes != "true" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" {
  2852. hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId))
  2853. }
  2854. */
  2855. resp, err := dockercli.ContainerCreate(ctx, config, hostConfig, networkingConfig, nil, containerName)
  2856. if err != nil {
  2857. if strings.Contains(err.Error(), "path does not exist") {
  2858. log.Printf("[ERROR] Not using permanent pipeline storage as storage folder %s does not exist. If you want permanent storage, create the %s folder then restart Orborus (1). Raw: %s", tenzirStorageFolder, tenzirStorageFolder, err)
  2859. skipPipelineMount = true
  2860. } else {
  2861. log.Printf("[ERROR] Failed to create Tenzir Node container: %v", err)
  2862. }
  2863. return err
  2864. }
  2865. if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" {
  2866. networkName := "shuffle_swarm_executions"
  2867. err = dockercli.NetworkConnect(ctx, networkName, resp.ID, nil)
  2868. if err != nil {
  2869. log.Printf("[ERROR] Error connecting tenzir container to network: %s", err)
  2870. }
  2871. }
  2872. err = dockercli.ContainerStart(ctx, containerName, containerStartOptions)
  2873. if err != nil {
  2874. if strings.Contains(err.Error(), "path does not exist") {
  2875. log.Printf("[ERROR] Not using permanent pipeline storage as storage folder %s does not exist. If you want permanent storage, create the %s folder then restart Orborus (2). Raw: %s", tenzirStorageFolder, tenzirStorageFolder, err)
  2876. skipPipelineMount = true
  2877. } else {
  2878. log.Printf("[ERROR] Failed to START Tenzir Node container: %v", err)
  2879. }
  2880. return err
  2881. }
  2882. log.Printf("[INFO] Tenzir Node container started successfully. Waiting for it to become available..")
  2883. time.Sleep(20 * time.Second)
  2884. err = checkTenzirNode()
  2885. if err != nil {
  2886. log.Printf("[ERROR] Tenzir connection not available: %s. IF the URL seems wrong, set SHUFFLE_PIPELINE_URL=http://<ip>:5160", err)
  2887. return err
  2888. }
  2889. log.Printf("[INFO] Successfully deployed Tenzir Node! Setting up default syslog listener on TCP/1514 AND UDP/1514")
  2890. command := `load_tcp "0.0.0.0:1514" { read_syslog } | import`
  2891. _, err = createPipeline(command, "default-syslog-tcp-514")
  2892. if err != nil {
  2893. log.Printf("[ERROR] Failed to create tcp syslog pipeline: %s", err)
  2894. return nil
  2895. }
  2896. command = `load_udp "0.0.0.0:1514", insert_newlines=true | read_syslog | import`
  2897. _, err = createPipeline(command, "default-syslog-udp-514")
  2898. if err != nil {
  2899. log.Printf("[ERROR] Failed to create udp syslog pipeline: %s", err)
  2900. return nil
  2901. }
  2902. return nil
  2903. }
  2904. func createNetworkIfNotExists(ctx context.Context, networkName, subnet, gateway string) error {
  2905. listOptions := network.ListOptions{}
  2906. networks, err := dockercli.NetworkList(ctx, listOptions)
  2907. if err != nil {
  2908. return err
  2909. }
  2910. for _, network := range networks {
  2911. if network.Name == networkName {
  2912. // Network exists
  2913. return nil
  2914. }
  2915. }
  2916. ipamConfig := &network.IPAM{
  2917. Config: []network.IPAMConfig{
  2918. {
  2919. Subnet: subnet,
  2920. Gateway: gateway,
  2921. },
  2922. },
  2923. }
  2924. networkCreate := network.CreateOptions{
  2925. //CheckDuplicate: true,
  2926. Driver: "bridge",
  2927. IPAM: ipamConfig,
  2928. }
  2929. _, err = dockercli.NetworkCreate(ctx, networkName, networkCreate)
  2930. if err != nil {
  2931. return err
  2932. }
  2933. return nil
  2934. }
  2935. func checkTenzirNode() error {
  2936. if tenzirDisabled && os.Getenv("SHUFFLE_SKIP_PIPELINES") == "true" && os.Getenv("SHUFFLE_PIPELINE_ENABLED") == "false" {
  2937. return errors.New("Pipelines are disabled by user with SHUFFLE_SKIP_PIPELINES (2)")
  2938. }
  2939. url := fmt.Sprintf("%s/api/v0/ping", pipelineUrl)
  2940. forwardMethod := "POST"
  2941. client := http.Client{
  2942. Timeout: 1 * time.Second,
  2943. }
  2944. req, err := http.NewRequest(forwardMethod, url, nil)
  2945. if err != nil {
  2946. log.Printf("[ERROR] Failed to create HTTP request: %s", err)
  2947. return err
  2948. }
  2949. resp, err := client.Do(req)
  2950. if err == nil && resp.StatusCode == http.StatusOK {
  2951. return nil
  2952. }
  2953. return fmt.Errorf("Tenzir node is not available due to: %s", err)
  2954. }
  2955. func createPipeline(command, identifier string) (string, error) {
  2956. //toBeDeleted := false
  2957. /*
  2958. // Pre-checked. No point here
  2959. pipelineId, err := searchPipeline(identifier)
  2960. if err != nil {
  2961. return "", err
  2962. }
  2963. */
  2964. url := fmt.Sprintf("%s/api/v0/pipeline/create", pipelineUrl)
  2965. forwardMethod := "POST"
  2966. /*
  2967. if err != nil {
  2968. if strings.Contains(fmt.Sprintf("%s", err), "no existing pipeline found") {
  2969. log.Printf("[INFO] No existing pipeline found with id: %s. Creating a new one!", identifier)
  2970. } else {
  2971. log.Printf("[ERROR] Failed to search for existing pipeline but continuing anyway : %s", err)
  2972. }
  2973. } else {
  2974. log.Printf("[INFO] an existing pipeline found with ID: %s. it will be deleted", pipelineId)
  2975. toBeDeleted = true
  2976. }
  2977. */
  2978. // if strings.Contains(command, "shuffler.io") {
  2979. // } else {
  2980. // var scheme string
  2981. // if strings.Contains(command, "http://") {
  2982. // scheme = "http://"
  2983. // } else if strings.Contains(command, "https://") {
  2984. // scheme = "https://"
  2985. // }
  2986. // startIndex := strings.Index(command, scheme)
  2987. // if startIndex != -1 {
  2988. // endIndex := startIndex + len(scheme)
  2989. // endIndex += strings.Index(command[endIndex:], "/")
  2990. // command = command[:startIndex] + baseUrl + command[endIndex:]
  2991. // }
  2992. // }
  2993. //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | sigma /var/lib/tenzir/rule.yaml"
  2994. //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | import"
  2995. // Make sure to escape them
  2996. //if strings.Contains(command, "/") {
  2997. // command = strings.ReplaceAll("\\\"", "", command)
  2998. // command = strings.ReplaceAll(command, "\"", "")
  2999. //}
  3000. requestBody := map[string]interface{}{
  3001. "definition": command,
  3002. "name": identifier,
  3003. "hidden": false,
  3004. "retry_delay": "500.0ms",
  3005. "unstoppable": true,
  3006. }
  3007. requestBodyJSON, err := json.Marshal(requestBody)
  3008. if err != nil {
  3009. log.Printf("[ERROR] failed marshalling body: %s", err)
  3010. return "", err
  3011. }
  3012. forwardData := bytes.NewBuffer(requestBodyJSON)
  3013. req, err := http.NewRequest(
  3014. forwardMethod,
  3015. url,
  3016. forwardData,
  3017. )
  3018. if err != nil {
  3019. log.Printf("[ERROR] Failed to create HTTP request: %s", err)
  3020. return "", err
  3021. }
  3022. req.Header.Set("Content-Type", "application/json")
  3023. client := &http.Client{Timeout: 10 * time.Second}
  3024. resp, err := client.Do(req)
  3025. if err != nil {
  3026. log.Printf("[ERROR] Failed to send HTTP request: %s", err)
  3027. return "", err
  3028. }
  3029. body, err := ioutil.ReadAll(resp.Body)
  3030. if err != nil {
  3031. log.Printf("[ERROR] Failed reading response body: %s", err)
  3032. return "", err
  3033. }
  3034. if strings.Contains(string(body), "error") {
  3035. log.Printf("[ERROR] Pipeline creation error resp (%d): %s", resp.StatusCode, string(body))
  3036. } else {
  3037. log.Printf("[DEBUG] Pipeline creation debug (%d): %s", resp.StatusCode, string(body))
  3038. }
  3039. defer resp.Body.Close()
  3040. if resp.StatusCode != 200 {
  3041. log.Printf("[DEBUG] status code is %d instead of 200", resp.StatusCode)
  3042. return "", fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
  3043. }
  3044. type PipelineResponse struct {
  3045. ID string `json:"id"`
  3046. Message string `json:"message"`
  3047. Severity string `json:"severity"`
  3048. }
  3049. var response PipelineResponse
  3050. if err := json.Unmarshal(body, &response); err != nil {
  3051. log.Printf("[ERROR] Failed unmarshalling response: %s", err)
  3052. return "", err
  3053. }
  3054. if response.ID == "" {
  3055. log.Printf("[ERROR] ID not found or empty in response. Severity: %#v, Message: %#v", response.Severity, response.Message)
  3056. return "", errors.New("Pipeline ID not found or empty in the response. See error logs.")
  3057. }
  3058. return response.ID, nil
  3059. }
  3060. func updatePipelineState(command, pipelineId, action string) (string, error) {
  3061. url := fmt.Sprintf("%s/api/v0/pipeline/update", pipelineUrl)
  3062. forwardMethod := "POST"
  3063. requestBody := map[string]interface{}{
  3064. "id": pipelineId,
  3065. "action": action,
  3066. /*
  3067. "autostart": map[string]bool{
  3068. "created": true,
  3069. "completed": false,
  3070. "failed": false,
  3071. },
  3072. "autodelete": map[string]bool{
  3073. "completed": false,
  3074. "failed": false,
  3075. "stopped": false,
  3076. },
  3077. */
  3078. }
  3079. requestBodyJSON, err := json.Marshal(requestBody)
  3080. if err != nil {
  3081. return "", err
  3082. }
  3083. log.Printf("[INFO] Updating pipeline %s with action %s to ensure it starts. Body: %s", pipelineId, action, string(requestBodyJSON))
  3084. forwardData := bytes.NewBuffer(requestBodyJSON)
  3085. req, err := http.NewRequest(
  3086. forwardMethod,
  3087. url,
  3088. forwardData,
  3089. )
  3090. if err != nil {
  3091. log.Printf("[ERROR] Failed to update HTTP request: %s", err)
  3092. return "", err
  3093. }
  3094. req.Header.Set("Content-Type", "application/json")
  3095. client := &http.Client{Timeout: 10 * time.Second}
  3096. resp, err := client.Do(req)
  3097. if err != nil {
  3098. log.Printf("[ERROR] Failed to send HTTP request: %s", err)
  3099. return "", err
  3100. }
  3101. defer resp.Body.Close()
  3102. if resp.StatusCode != http.StatusOK {
  3103. return "", fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
  3104. }
  3105. body, err := ioutil.ReadAll(resp.Body)
  3106. if err != nil {
  3107. return "", err
  3108. }
  3109. var responseData struct {
  3110. Pipeline struct {
  3111. State string `json:"state"`
  3112. } `json:"pipeline"`
  3113. }
  3114. if err := json.Unmarshal(body, &responseData); err != nil {
  3115. return "", err
  3116. }
  3117. return responseData.Pipeline.State, nil
  3118. }
  3119. func deletePipeline(pipelineId string) error {
  3120. requestBody := map[string]string{
  3121. "id": pipelineId,
  3122. }
  3123. url := fmt.Sprintf("%s/api/v0/pipeline/delete", pipelineUrl)
  3124. forwardMethod := "POST"
  3125. requestBodyJSON, err := json.Marshal(requestBody)
  3126. if err != nil {
  3127. log.Println("[ERROR] failed marshalling request body:", err)
  3128. return err
  3129. }
  3130. forwardData := bytes.NewBuffer(requestBodyJSON)
  3131. req, err := http.NewRequest(
  3132. forwardMethod,
  3133. url,
  3134. forwardData,
  3135. )
  3136. if err != nil {
  3137. log.Printf("[ERROR] Failed to delete HTTP request: %s", err)
  3138. return err
  3139. }
  3140. req.Header.Set("Content-Type", "application/json")
  3141. client := &http.Client{Timeout: 10 * time.Second}
  3142. resp, err := client.Do(req)
  3143. if err != nil {
  3144. log.Printf("[ERROR] Failed to send HTTP request: %s", err)
  3145. return err
  3146. }
  3147. defer resp.Body.Close()
  3148. if resp.StatusCode != 200 {
  3149. log.Printf("[DEBUG] The deletion of pipeline with ID: %s is unsucessful as status code is NOT 200 !!!", pipelineId)
  3150. return fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
  3151. }
  3152. log.Printf("[INFO] Pipeline with ID: %s deleted successfully", pipelineId)
  3153. pipelines = []shuffle.PipelineInfo{}
  3154. return nil
  3155. }
  3156. // Lists the pipelines from the API exactly as they are. Definition is set up in Shuffle structs
  3157. func listPipelines() ([]shuffle.PipelineInfo, error) {
  3158. responseData := shuffle.PipelineInfoWrapper{}
  3159. if tenzirDisabled {
  3160. return responseData.Pipelines, errors.New("Tenzir is disabled")
  3161. }
  3162. var reqBody []byte
  3163. url := fmt.Sprintf("%s/api/v0/pipeline/list", pipelineUrl)
  3164. client := http.Client{
  3165. Timeout: 2 * time.Second,
  3166. }
  3167. req, err := http.NewRequest(
  3168. "POST",
  3169. url,
  3170. bytes.NewBuffer(reqBody),
  3171. )
  3172. if err != nil {
  3173. return responseData.Pipelines, err
  3174. }
  3175. req.Header.Set("Content-Type", "application/json")
  3176. resp, err := client.Do(req)
  3177. if err != nil {
  3178. return responseData.Pipelines, err
  3179. }
  3180. defer resp.Body.Close()
  3181. if resp.StatusCode != http.StatusOK {
  3182. return responseData.Pipelines, fmt.Errorf("Got the status code %d instead of 200 from Pipeline node", resp.StatusCode)
  3183. }
  3184. body, err := ioutil.ReadAll(resp.Body)
  3185. if err != nil {
  3186. return responseData.Pipelines, err
  3187. }
  3188. if err := json.Unmarshal(body, &responseData); err != nil {
  3189. return responseData.Pipelines, err
  3190. }
  3191. return responseData.Pipelines, nil
  3192. }
  3193. func searchPipeline(identifier string) (string, error) {
  3194. allPipelines, err := listPipelines()
  3195. if err != nil {
  3196. return "", err
  3197. }
  3198. for _, pipeline := range allPipelines {
  3199. if pipeline.Name == identifier {
  3200. return pipeline.ID, nil
  3201. }
  3202. }
  3203. return "", errors.New("no existing pipeline found with name")
  3204. }
  3205. func handleFileCategoryChange(ruleType string) error {
  3206. apiEndpoint := fmt.Sprintf("%s/api/v1/files/namespaces/%s", baseUrl, ruleType)
  3207. req, err := http.NewRequest("GET", apiEndpoint, nil)
  3208. if err != nil {
  3209. return err
  3210. }
  3211. if len(pipelineApikey) == 0 {
  3212. //var auth = os.Getenv("AUTH")
  3213. //var org = os.Getenv("ORG")
  3214. if len(auth) > 0 && len(org) > 0 {
  3215. pipelineApikey = auth
  3216. } else {
  3217. return errors.New("Shuffle API-key not set for Pipelines: SHUFFLE_PIPELINE_AUTH=<apikey>")
  3218. }
  3219. }
  3220. req.Header.Add("Authorization", "Bearer "+pipelineApikey)
  3221. if len(org) > 0 {
  3222. req.Header.Add("Org-Id", org)
  3223. }
  3224. client := shuffle.GetExternalClient(apiEndpoint)
  3225. resp, err := client.Do(req)
  3226. if err != nil {
  3227. return err
  3228. }
  3229. defer resp.Body.Close()
  3230. if resp.StatusCode != http.StatusOK {
  3231. return fmt.Errorf("Received non-200 response '%d' from backend URL %s. ", resp.StatusCode, apiEndpoint)
  3232. }
  3233. out, err := os.Create("files.zip")
  3234. if err != nil {
  3235. return err
  3236. }
  3237. defer out.Close()
  3238. defer os.Remove("files.zip")
  3239. _, err = io.Copy(out, resp.Body)
  3240. if err != nil {
  3241. log.Printf("[ERROR] Failed to io.Copy ZIP file content: %s", err)
  3242. return err
  3243. }
  3244. //log.Println("[DEBUG] ZIP file downloaded successfully.")
  3245. tenzirStorageFolder := os.Getenv("SHUFFLE_STORAGE_FOLDER")
  3246. if len(tenzirStorageFolder) == 0 {
  3247. tenzirStorageFolder = "/tmp/"
  3248. }
  3249. tenzirStorageFolder = strings.TrimRight(tenzirStorageFolder, "/")
  3250. sigmaPath := fmt.Sprintf("%s/%s_rules", tenzirStorageFolder, ruleType)
  3251. err = extractZIP("files.zip", sigmaPath)
  3252. if err != nil {
  3253. log.Printf("[ERROR] Failed to extract ZIP file: %s", err)
  3254. return err
  3255. }
  3256. log.Printf("[DEBUG] Detection files copied to '%s' successfully.", sigmaPath)
  3257. return nil
  3258. }
  3259. func extractZIP(zipFile, destDir string) error {
  3260. r, err := zip.OpenReader(zipFile)
  3261. if err != nil {
  3262. return err
  3263. }
  3264. // FInd size of the zip
  3265. var totalSize uint64
  3266. for _, f := range r.File {
  3267. totalSize += f.UncompressedSize64
  3268. }
  3269. log.Printf("[DEBUG] Total size of the ZIP file: %d bytes", totalSize)
  3270. defer r.Close()
  3271. if err := os.MkdirAll(destDir, 0755); err != nil {
  3272. return err
  3273. }
  3274. log.Printf("[DEBUG] Total files to extract: %d", len(r.File))
  3275. for _, f := range r.File {
  3276. // Fix path traversal
  3277. if strings.Contains(f.Name, "..") {
  3278. return fmt.Errorf("illegal file name: %s", f.Name)
  3279. }
  3280. err := extractFile(f, destDir)
  3281. if err != nil {
  3282. return err
  3283. }
  3284. }
  3285. return nil
  3286. }
  3287. func extractFile(f *zip.File, destDir string) error {
  3288. rc, err := f.Open()
  3289. if err != nil {
  3290. return err
  3291. }
  3292. defer rc.Close()
  3293. path := filepath.Join(destDir, f.Name)
  3294. out, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, f.Mode())
  3295. if err != nil {
  3296. return err
  3297. }
  3298. defer out.Close()
  3299. _, err = io.Copy(out, rc)
  3300. return err
  3301. }
  3302. func copyToTenzir(srcPath, destPath string) error {
  3303. containerName := "tenzir-node"
  3304. checkCmd := exec.Command("docker", "exec", containerName, "test", "-d", destPath)
  3305. if err := checkCmd.Run(); err == nil {
  3306. rmCmd := exec.Command("docker", "exec", "-u", "root", containerName, "rm", "-rf", destPath)
  3307. if err := rmCmd.Run(); err != nil {
  3308. return fmt.Errorf("error removing existing directory in container: %v", err)
  3309. }
  3310. }
  3311. cpCmd := exec.Command("docker", "cp", srcPath, fmt.Sprintf("%s:%s", containerName, destPath))
  3312. var out bytes.Buffer
  3313. cpCmd.Stdout = &out
  3314. cpCmd.Stderr = &out
  3315. err := cpCmd.Run()
  3316. if err != nil {
  3317. return fmt.Errorf("error copying files: %v, output: %s", err, out.String())
  3318. }
  3319. return nil
  3320. }
  3321. func removeFileCategory(ruleType string) error {
  3322. tenzirStorageFolder := os.Getenv("SHUFFLE_STORAGE_FOLDER")
  3323. if len(tenzirStorageFolder) == 0 {
  3324. tenzirStorageFolder = "/tmp/"
  3325. }
  3326. tenzirStorageFolder = strings.TrimRight(tenzirStorageFolder, "/")
  3327. //sigmaPath := "/var/lib/tenzir/sigma_rules/*"
  3328. rulePath := fmt.Sprintf("%s/%s_rules", tenzirStorageFolder, ruleType)
  3329. err := os.RemoveAll(rulePath)
  3330. if err != nil {
  3331. return fmt.Errorf("Error removing category files in %s: %v", rulePath, err)
  3332. }
  3333. log.Printf("[INFO] Removed all local category data in %s", rulePath)
  3334. return nil
  3335. }
  3336. // curl https://get.tenzir.app | sh
  3337. func removeFile(fileName string) error {
  3338. containerName := "tenzir-node"
  3339. srcPath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", fileName)
  3340. checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
  3341. if err := checkSrcCmd.Run(); err != nil {
  3342. // If the file does not exist, simply return nil
  3343. if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
  3344. log.Printf("[ERROR] No such file: %s, nothing to delete\n", srcPath)
  3345. return nil
  3346. }
  3347. return fmt.Errorf("error checking source file: %v", err)
  3348. }
  3349. return removePath(containerName, srcPath)
  3350. }
  3351. func removePath(containerName, path string) error {
  3352. rmCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("rm -rf %s", path))
  3353. output, err := rmCmd.CombinedOutput()
  3354. if err != nil {
  3355. return fmt.Errorf("error removing path: %v, output: %s", err, output)
  3356. }
  3357. return nil
  3358. }
  3359. func sendPipelineHealthStatus() (shuffle.LakeConfig, error) {
  3360. pipelinePayload := shuffle.LakeConfig{
  3361. Enabled: false,
  3362. Pipelines: []shuffle.PipelineInfo{},
  3363. }
  3364. if tenzirDisabled {
  3365. return pipelinePayload, nil
  3366. }
  3367. // To not spam down the list API too much
  3368. randint := rand.Intn(5)
  3369. if len(pipelines) == 0 || randint == 0 {
  3370. pipelineDef, err := listPipelines()
  3371. if err == nil || len(pipelines) > 0 {
  3372. pipelines = pipelineDef
  3373. pipelinePayload.Pipelines = pipelines
  3374. }
  3375. } else {
  3376. pipelinePayload.Pipelines = pipelines
  3377. }
  3378. err := deployTenzirNode()
  3379. if err != nil {
  3380. if (!strings.Contains(err.Error(), "SHUFFLE_SKIP_PIPELINES") && !strings.Contains(err.Error(), "Kubernetes not implemented for Tenzir node")) && !strings.Contains(err.Error(), "Tenzir Node is already running") && !strings.Contains(err.Error(), "docker daemon") {
  3381. log.Printf("[ERROR] Tenzir node connection problem: %s", err)
  3382. } else {
  3383. //tenzirDisabled = true
  3384. if debug {
  3385. log.Printf("[WARNING] Disabling pipelines: %s. You will need to restart the Orborus to fix this.", err)
  3386. }
  3387. }
  3388. return pipelinePayload, err
  3389. }
  3390. pipelinePayload.Enabled = true
  3391. // No direct sending.
  3392. return pipelinePayload, nil
  3393. }
  3394. func disableRule(fileName string) error {
  3395. containerName := "tenzir-node"
  3396. srcPath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", fileName)
  3397. destDir := "/var/lib/tenzir/disabled_rules"
  3398. destPath := fmt.Sprintf("%s/%s", destDir, fileName)
  3399. checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
  3400. if err := checkSrcCmd.Run(); err != nil {
  3401. if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
  3402. fmt.Printf("File does not exist: %s\n", srcPath)
  3403. return nil // Nothing to disable
  3404. }
  3405. return fmt.Errorf("error checking source file: %v", err)
  3406. }
  3407. checkDestDirCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mkdir -p %s", destDir))
  3408. if err := checkDestDirCmd.Run(); err != nil {
  3409. return fmt.Errorf("error ensuring destination directory exists: %v", err)
  3410. }
  3411. moveCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mv %s %s", srcPath, destPath))
  3412. if err := moveCmd.Run(); err != nil {
  3413. return fmt.Errorf("error moving file: %v", err)
  3414. }
  3415. fmt.Printf("File %s moved to %s successfully.\n", fileName, destDir)
  3416. return nil
  3417. }
  3418. func enableRule(fileName string) error {
  3419. containerName := "tenzir-node"
  3420. srcPath := fmt.Sprintf("/var/lib/tenzir/disabled_rules/%s", fileName)
  3421. destDir := "/var/lib/tenzir/sigma_rules"
  3422. destPath := fmt.Sprintf("%s/%s", destDir, fileName)
  3423. checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
  3424. if err := checkSrcCmd.Run(); err != nil {
  3425. if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
  3426. fmt.Printf("File does not exist: %s\n", srcPath)
  3427. return nil // Nothing to enable
  3428. }
  3429. return fmt.Errorf("error checking source file: %v", err)
  3430. }
  3431. checkDestDirCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mkdir -p %s", destDir))
  3432. if err := checkDestDirCmd.Run(); err != nil {
  3433. return fmt.Errorf("error ensuring destination directory exists: %v", err)
  3434. }
  3435. moveCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mv %s %s", srcPath, destPath))
  3436. if err := moveCmd.Run(); err != nil {
  3437. return fmt.Errorf("error moving file: %v", err)
  3438. }
  3439. fmt.Printf("[DEBUG] File %s moved to %s successfully.\n", fileName, destDir)
  3440. return nil
  3441. }
  3442. // Is this ok to do with Docker? idk :)
  3443. func getRunningWorkers(ctx context.Context, workerTimeout int) int {
  3444. //log.Printf("[DEBUG] Getting running workers with API version %s", dockerApiVersion)
  3445. counter := 0
  3446. if isKubernetes == "true" {
  3447. log.Printf("[INFO] Getting running workers in kubernetes")
  3448. thresholdTime := time.Now().Add(time.Duration(-workerTimeout) * time.Second)
  3449. clientset, _, err := shuffle.GetKubernetesClient()
  3450. if err != nil {
  3451. log.Printf("[ERROR] Failed getting kubernetes client: %s", err)
  3452. return 0
  3453. }
  3454. pods, podErr := clientset.CoreV1().Pods(kubernetesNamespace).List(ctx, metav1.ListOptions{
  3455. LabelSelector: "app.kubernetes.io/name=shuffle-worker",
  3456. })
  3457. if podErr != nil {
  3458. log.Printf("[ERROR] Failed getting running workers: %s", podErr)
  3459. return 0
  3460. }
  3461. for _, pod := range pods.Items {
  3462. if pod.Status.Phase == "Running" && pod.CreationTimestamp.Time.After(thresholdTime) {
  3463. counter++
  3464. }
  3465. }
  3466. if counter > 0 {
  3467. log.Printf("[INFO] Found %d running workers in Orborus", counter)
  3468. }
  3469. } else {
  3470. containers, err := dockercli.ContainerList(ctx, container.ListOptions{
  3471. All: true,
  3472. })
  3473. // Automatically updates the version
  3474. if err != nil {
  3475. log.Printf("[ERROR] Error getting containers from Docker: %s", err)
  3476. newVersionSplit := strings.Split(fmt.Sprintf("%s", err), "version is")
  3477. if len(newVersionSplit) > 1 {
  3478. //dockerApiVersion = strings.TrimSpace(newVersionSplit[1])
  3479. log.Printf("[DEBUG] WANT to change the API version to default to %s?", strings.TrimSpace(newVersionSplit[1]))
  3480. }
  3481. return maxConcurrency
  3482. }
  3483. currenttime := time.Now().Unix()
  3484. for _, container := range containers {
  3485. // Skip random containers. Only handle things related to Shuffle.
  3486. if !strings.Contains(container.Image, baseimagename) {
  3487. shuffleFound := false
  3488. for _, item := range container.Labels {
  3489. if item == "shuffle" {
  3490. shuffleFound = true
  3491. break
  3492. }
  3493. }
  3494. // Check image name
  3495. if !shuffleFound {
  3496. continue
  3497. }
  3498. //} else {
  3499. // log.Printf("NAME: %s", container.Image)
  3500. }
  3501. for _, name := range container.Names {
  3502. // FIXME - add name_version_uid_uid regex check as well
  3503. if !strings.HasPrefix(name, "/worker") {
  3504. continue
  3505. }
  3506. //log.Printf("Time: %d - %d", currenttime-container.Created, int64(workerTimeout))
  3507. if container.State == "running" && currenttime-container.Created < int64(workerTimeout) {
  3508. counter += 1
  3509. break
  3510. }
  3511. }
  3512. }
  3513. }
  3514. return counter
  3515. }
  3516. // FIXME - add this to remove exited workers
  3517. // Should it check what happened to the execution? idk
  3518. func zombiecheck(ctx context.Context, workerTimeout int) error {
  3519. isK8s := isKubernetes == "true"
  3520. executionIds = []string{}
  3521. if swarmConfig == "run" || swarmConfig == "swarm" || isK8s {
  3522. //log.Printf("[DEBUG] Skipping Zombie check due to new execution model (swarm)")
  3523. return nil
  3524. }
  3525. log.Println("[INFO] Looking for old containers to remove")
  3526. containers, err := dockercli.ContainerList(ctx, container.ListOptions{
  3527. All: true,
  3528. })
  3529. if err != nil {
  3530. log.Printf("[ERROR] Failed creating Containerlist: %s", err)
  3531. return err
  3532. }
  3533. containerNames := map[string]string{}
  3534. stopContainers := []string{}
  3535. removeContainers := []string{}
  3536. log.Printf("[INFO] Baseimage: %s, Workertimeout: %d", baseimagename, int64(workerTimeout))
  3537. //baseString := `/bin/sh -c 'python app.py --log-level DEBUG'`
  3538. baseString := `python app.py`
  3539. for _, container := range containers {
  3540. // Skip random containers. Only handle things related to Shuffle.
  3541. if !strings.Contains(container.Image, baseimagename) && !strings.Contains(container.Command, baseString) && !strings.Contains(container.Command, "walkoff") && container.Command != "./worker" {
  3542. shuffleFound := false
  3543. for _, item := range container.Labels {
  3544. if item == "shuffle" {
  3545. shuffleFound = true
  3546. break
  3547. }
  3548. }
  3549. // Check image name
  3550. if !shuffleFound {
  3551. //log.Printf("[DEBUG] Zombie container skip: %#v, %s", container.Labels, container.Image)
  3552. continue
  3553. }
  3554. //} else {
  3555. // log.Printf("NAME: %s", container.Image)
  3556. } else {
  3557. //log.Printf("Img: %s", container.Image)
  3558. //log.Printf("Names: %s", container.Names)
  3559. }
  3560. for _, name := range container.Names {
  3561. // FIXME - add name_version_uid_uid regex check as well
  3562. if strings.HasPrefix(name, "/shuffle") && !strings.HasPrefix(name, "/shuffle-subflow") {
  3563. continue
  3564. }
  3565. currenttime := time.Now().Unix()
  3566. //log.Printf("[INFO] (%s) NAME: %s. TIME: %d", container.State, name, currenttime-container.Created)
  3567. // Need to check time here too because a container can be removed the same instant as its created
  3568. if container.State != "running" && currenttime-container.Created > int64(workerTimeout) {
  3569. removeContainers = append(removeContainers, container.ID)
  3570. containerNames[container.ID] = name
  3571. }
  3572. // stopcontainer & removecontainer
  3573. //log.Printf("Time: %d - %d", currenttime-container.Created, int64(workerTimeout))
  3574. if container.State == "running" && currenttime-container.Created > int64(workerTimeout) {
  3575. stopContainers = append(stopContainers, container.ID)
  3576. containerNames[container.ID] = name
  3577. }
  3578. }
  3579. }
  3580. // FIXME - add killing of apps with same execution ID too
  3581. log.Printf("[INFO] Should STOP and remove %d containers.", len(stopContainers))
  3582. var options container.StopOptions
  3583. for _, containername := range stopContainers {
  3584. log.Printf("[INFO] Stopping and removing container %s", containerNames[containername])
  3585. go dockercli.ContainerStop(ctx, containername, options)
  3586. removeContainers = append(removeContainers, containername)
  3587. }
  3588. removeOptions := container.RemoveOptions{
  3589. RemoveVolumes: true,
  3590. Force: true,
  3591. }
  3592. log.Printf("[INFO] Should REMOVE %d containers.", len(removeContainers))
  3593. for _, containername := range removeContainers {
  3594. dockercli.ContainerRemove(ctx, containername, removeOptions)
  3595. }
  3596. return nil
  3597. }
  3598. func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest, image string, env []string) error {
  3599. parsedRequest := shuffle.OrborusExecutionRequest{
  3600. ExecutionId: workflowExecution.ExecutionId,
  3601. Authorization: workflowExecution.Authorization,
  3602. BaseUrl: os.Getenv("BASE_URL"),
  3603. EnvironmentName: os.Getenv("ENVIRONMENT_NAME"),
  3604. Timezone: os.Getenv("TZ"),
  3605. Cleanup: os.Getenv("CLEANUP"),
  3606. HTTPProxy: os.Getenv("HTTP_PROXY"),
  3607. HTTPSProxy: os.Getenv("HTTPS_PROXY"),
  3608. ShufflePassProxyToApp: os.Getenv("SHUFFLE_PASS_APP_PROXY"),
  3609. WorkerServerUrl: os.Getenv("SHUFFLE_WORKER_SERVER_URL"),
  3610. }
  3611. parsedBaseurl := baseUrl
  3612. if strings.Contains(baseUrl, ":") {
  3613. baseUrlSplit := strings.Split(baseUrl, ":")
  3614. if len(baseUrlSplit) >= 3 {
  3615. parsedBaseurl = strings.Join(baseUrlSplit[0:2], ":")
  3616. }
  3617. }
  3618. data, err := json.Marshal(parsedRequest)
  3619. if err != nil {
  3620. log.Printf("[ERROR] Failed marshalling worker request: %s", err)
  3621. return err
  3622. }
  3623. streamUrl := fmt.Sprintf("http://shuffle-workers:33333/api/v1/execute")
  3624. if containerId == "" || containerId == "shuffle-orborus" {
  3625. streamUrl = fmt.Sprintf("%s:33333/api/v1/execute", parsedBaseurl)
  3626. }
  3627. identifier := "shuffle-workers"
  3628. if isKubernetes == "true" {
  3629. // FIXME: Do we need this to map the cluster?
  3630. //if shuffle.IsRunningInCluster() {
  3631. //log.Printf("[INFO] Running in Kubernetes cluster")
  3632. // try getting the k8s worker server url
  3633. //}
  3634. }
  3635. if strings.Contains(streamUrl, "shuffler.io") || strings.Contains(streamUrl, "localhost") || strings.Contains(streamUrl, "127.0.0.1") || strings.Contains(streamUrl, "shuffle-backend") {
  3636. // Specific to debugging
  3637. if len(workerServerUrl) == 0 {
  3638. if debug {
  3639. log.Printf("[INFO] Using default worker server url as previous is invalid: %s. Swapping to shuffle-workers:33333", streamUrl)
  3640. }
  3641. }
  3642. streamUrl = fmt.Sprintf("http://shuffle-workers:33333/api/v1/execute")
  3643. }
  3644. if len(workerServerUrl) > 0 {
  3645. // Check if a port is supplied or not
  3646. if strings.Contains(workerServerUrl, "/api/v1/execute") {
  3647. streamUrl = workerServerUrl
  3648. } else {
  3649. streamUrl = fmt.Sprintf("%s/api/v1/execute", workerServerUrl)
  3650. if !strings.Contains(workerServerUrl, ":") {
  3651. streamUrl = fmt.Sprintf("%s:33333/api/v1/execute", workerServerUrl)
  3652. }
  3653. }
  3654. }
  3655. client := &http.Client{
  3656. //Transport: &http.Transport{
  3657. // TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
  3658. //},
  3659. Timeout: time.Duration(120 * time.Second),
  3660. }
  3661. if debug {
  3662. log.Printf("[DEBUG][%s] Worker request to be sent to URL: %s", workflowExecution.ExecutionId, streamUrl)
  3663. }
  3664. req, err := http.NewRequest(
  3665. "POST",
  3666. streamUrl,
  3667. bytes.NewBuffer([]byte(data)),
  3668. )
  3669. if err != nil {
  3670. log.Printf("[ERROR] Failed creating worker request: %s", err)
  3671. if strings.Contains(fmt.Sprintf("%s", err), "connection refused") || strings.Contains(fmt.Sprintf("%s", err), "EOF") {
  3672. workerImage := fmt.Sprintf("ghcr.io/shuffle/shuffle-worker:%s", workerVersion)
  3673. if len(newWorkerImage) > 0 {
  3674. workerImage = newWorkerImage
  3675. }
  3676. if isKubernetes == "true" {
  3677. deployK8sWorker(workerImage, identifier, env)
  3678. } else {
  3679. deployServiceWorkers(workerImage)
  3680. }
  3681. time.Sleep(time.Duration(10) * time.Second)
  3682. //err = sendWorkerRequest(executionRequest)
  3683. }
  3684. return err
  3685. }
  3686. newresp, err := client.Do(req)
  3687. if err != nil {
  3688. // Connection refused?
  3689. if !strings.Contains(fmt.Sprintf("%s", err), "timeout") {
  3690. log.Printf("[ERROR][%s] Error running worker request to %s (1): %s", workflowExecution.ExecutionId, streamUrl, err)
  3691. }
  3692. if strings.Contains(fmt.Sprintf("%s", err), "connection refused") || strings.Contains(fmt.Sprintf("%s", err), "EOF") {
  3693. workerImage := fmt.Sprintf("ghcr.io/shuffle/shuffle-worker:%s", workerVersion)
  3694. if len(newWorkerImage) > 0 {
  3695. workerImage = newWorkerImage
  3696. }
  3697. if isKubernetes == "true" {
  3698. deployK8sWorker(workerImage, identifier, env)
  3699. } else {
  3700. deployServiceWorkers(workerImage)
  3701. }
  3702. time.Sleep(time.Duration(10) * time.Second)
  3703. //err = sendWorkerRequest(executionRequest)
  3704. }
  3705. return err
  3706. }
  3707. defer newresp.Body.Close()
  3708. body, err := ioutil.ReadAll(newresp.Body)
  3709. if err != nil {
  3710. log.Printf("[ERROR] Failed reading body in worker request body to worker on %s: %s", streamUrl, err)
  3711. return err
  3712. }
  3713. window.AddEvent(time.Now())
  3714. if newresp.StatusCode != 200 {
  3715. log.Printf("[WARNING] POTENTIAL error running worker request (2) - status code is %d for %s, not 200. Body: %s", newresp.StatusCode, streamUrl, string(body))
  3716. // In case of old executions
  3717. if strings.Contains(strings.ToLower(string(body)), "bad status ") {
  3718. return nil
  3719. }
  3720. if strings.Contains(strings.ToLower(string(body)), "no apps to handle") {
  3721. return nil
  3722. }
  3723. return errors.New(fmt.Sprintf("Bad statuscode from worker: %d - expecting 200", newresp.StatusCode))
  3724. }
  3725. _ = body
  3726. debugCommand := fmt.Sprintf("docker service logs shuffle-workers 2>&1 -f | grep %s", workflowExecution.ExecutionId)
  3727. if isKubernetes == "true" {
  3728. debugCommand = fmt.Sprintf("kubectl logs -n %s deployment/shuffle-workers | grep %s", kubernetesNamespace, workflowExecution.ExecutionId)
  3729. }
  3730. log.Printf("[DEBUG][%s] Ran worker from requests. Worker URL: %s. DEBUGGING:\n%s", workflowExecution.ExecutionId, streamUrl, debugCommand)
  3731. return nil
  3732. }
  3733. // 0x0elliot:
  3734. // let's never increase worker replicas.
  3735. // in our tests, workers replicas mattered a lot less.
  3736. // edge-case: subflows are helped with when worker replicas are higher.
  3737. func AutoScale(ctx context.Context) {
  3738. if os.Getenv("SHUFFLE_SCALE_REPLICAS") != "" {
  3739. return
  3740. }
  3741. ticker := time.NewTicker(1 * time.Second)
  3742. coolDownPeriod := 10 * time.Second
  3743. queuePerMinuteInt = 20
  3744. if os.Getenv("SHUFFLE_QUEUE_PER_MINUTE") != "" {
  3745. var err error
  3746. queuePerMinuteInt, err = strconv.Atoi(os.Getenv("SHUFFLE_QUEUE_PER_MINUTE"))
  3747. if err != nil {
  3748. log.Printf("[WARNING] Cannot convert %s to int. Using default value for it: %d", queuePerMinute, queuePerMinuteInt)
  3749. }
  3750. }
  3751. lastScaleTime := time.Now()
  3752. currentWorkers := currentWokerCount(ctx, dockercli)
  3753. for {
  3754. select {
  3755. case <-ctx.Done():
  3756. return
  3757. case <-ticker.C:
  3758. if time.Since(lastScaleTime) < (coolDownPeriod) {
  3759. continue
  3760. }
  3761. currentRequestCount := window.CountEvents(time.Now())
  3762. requiredReplicas := 0
  3763. if currentRequestCount >= queuePerMinuteInt*currentWorkers {
  3764. // FIXME: Hardcoded Max Replicas should be 6
  3765. requiredReplicas = int(math.Min(float64(6), float64(currentRequestCount/queuePerMinuteInt)+1))
  3766. }
  3767. if requiredReplicas > 0 {
  3768. err := scaleService(ctx, dockercli, uint64(requiredReplicas))
  3769. if err != nil {
  3770. log.Printf("[ERROR] Failed to scale service: %s", err)
  3771. } else {
  3772. lastScaleTime = time.Now()
  3773. currentWorkers = currentWokerCount(ctx, dockercli)
  3774. }
  3775. }
  3776. }
  3777. }
  3778. }
  3779. func scaleService(ctx context.Context, client *dockerclient.Client, replicas uint64) error {
  3780. service, _, err := client.ServiceInspectWithRaw(ctx, "shuffle-workers", types.ServiceInspectOptions{})
  3781. if err != nil {
  3782. return err
  3783. }
  3784. if service.Spec.Mode.Replicated == nil {
  3785. return errors.New("Service cannot be replicated")
  3786. }
  3787. if *service.Spec.Mode.Replicated.Replicas >= replicas {
  3788. return nil
  3789. }
  3790. service.Spec.Mode.Replicated.Replicas = &replicas
  3791. _, err = dockercli.ServiceUpdate(ctx, service.ID, service.Version, service.Spec, types.ServiceUpdateOptions{})
  3792. if err != nil {
  3793. return err
  3794. }
  3795. log.Printf("[INFO] Scaled shuffle-worker to %d replicas", replicas)
  3796. return nil
  3797. }
  3798. func currentWokerCount(ctx context.Context, client *dockerclient.Client) int {
  3799. service, _, err := client.ServiceInspectWithRaw(ctx, "shuffle-workers", types.ServiceInspectOptions{})
  3800. if err != nil {
  3801. return 0
  3802. }
  3803. if service.Spec.Mode.Replicated == nil {
  3804. return 0
  3805. }
  3806. return int(*service.Spec.Mode.Replicated.Replicas)
  3807. }
  3808. func queueScaleFactor(numQueue int, queuePerMin int) float64 {
  3809. if numQueue > queuePerMin {
  3810. queuePressure := float64(numQueue) / float64(queuePerMin)
  3811. return 1.0 + math.Min(queuePressure-1.0, 1.0)
  3812. }
  3813. return 1.0
  3814. }
  3815. func checkMemcached(ctx context.Context, dockercli *dockerclient.Client) (bool, error) {
  3816. containerName := "shuffle-cache"
  3817. continer, err := dockercli.ContainerInspect(context.Background(), containerName)
  3818. if err != nil {
  3819. if dockerclient.IsErrNotFound(err) {
  3820. return false, nil
  3821. }
  3822. return false, err
  3823. }
  3824. networkName := "shuffle_swarm_executions"
  3825. err = dockercli.NetworkConnect(ctx, networkName, containerName, nil)
  3826. if err != nil {
  3827. log.Printf("[WARNING] Failed connecting memcached container to network: %s", err)
  3828. }
  3829. if continer.State.Running == false {
  3830. log.Printf("[INFO] Container %s exists but is not running. Attempting to start it.", containerName)
  3831. err = dockercli.ContainerStart(ctx, containerName, container.StartOptions{})
  3832. if err != nil {
  3833. log.Printf("[ERROR] Failed to start container %s: %v", containerName, err)
  3834. return false, err
  3835. }
  3836. log.Printf("[INFO] Successfully started container %s.", containerName)
  3837. return true, nil
  3838. }
  3839. return continer.State.Running, nil
  3840. }
  3841. func deployMemcached(dockercli *dockerclient.Client) error {
  3842. if os.Getenv("SHUFFLE_MEMCACHED") != "" {
  3843. return errors.New("Memcached already running")
  3844. }
  3845. defaultMem := "1024"
  3846. log.Printf("[INFO] Spanning a default memcached container to handle the distribution between cache across different workers. Default memory assigned %s", defaultMem)
  3847. ctx := context.Background()
  3848. memcachedImage := "docker.io/library/memcached:latest"
  3849. containerConfig := &container.Config{
  3850. Image: memcachedImage,
  3851. Cmd: []string{"-m", defaultMem},
  3852. }
  3853. hostConfig := &container.HostConfig{
  3854. PortBindings: nat.PortMap{
  3855. "11211/tcp": []nat.PortBinding{{HostPort: "11211"}},
  3856. },
  3857. }
  3858. _, _, err := dockercli.ImageInspectWithRaw(ctx, memcachedImage)
  3859. if dockerclient.IsErrNotFound(err) {
  3860. log.Printf("[DEBUG] Pulling image %s. This may take a while.", memcachedImage)
  3861. pullOptions := image.PullOptions{}
  3862. out, err := dockercli.ImagePull(ctx, memcachedImage, pullOptions)
  3863. if err != nil {
  3864. log.Printf("[ERROR] Failed to pull the memcached image: %s", err)
  3865. return err
  3866. }
  3867. defer out.Close()
  3868. io.Copy(io.Discard, out)
  3869. } else if err != nil {
  3870. return err
  3871. }
  3872. containerName := "shuffle-cache"
  3873. resp, err := dockercli.ContainerCreate(ctx, containerConfig, hostConfig, nil, nil, containerName)
  3874. if err != nil {
  3875. log.Printf("[ERROR] Error spanning memcached continer: %s", err)
  3876. return err
  3877. }
  3878. if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" {
  3879. networkName := "shuffle_swarm_executions"
  3880. err = dockercli.NetworkConnect(ctx, networkName, resp.ID, nil)
  3881. if err != nil {
  3882. log.Printf("[ERROR] Error connecting tenzir container to network: %s", err)
  3883. }
  3884. }
  3885. err = dockercli.ContainerStart(ctx, resp.ID, container.StartOptions{})
  3886. if err != nil {
  3887. log.Printf("[ERROR] Error starting memcached continer: %s", err)
  3888. return err
  3889. }
  3890. networkName := "shuffle_swarm_executions"
  3891. err = dockercli.NetworkConnect(ctx, networkName, resp.ID, nil)
  3892. if err != nil {
  3893. log.Printf("[ERROR] Error connecting memcached container to network: %s", err)
  3894. }
  3895. log.Printf("[INFO] Memcached container started successfully at port 11211")
  3896. return nil
  3897. }
  3898. // How do we get the cpu usage? maybe just get the number of requests (much more useful for apps)
  3899. /*
  3900. func nodesResourceUsage(ctx context.Context, client *dockerclient.Client) error {
  3901. nodes, err := client.NodeList(ctx, types.NodeListOptions{})
  3902. if err != nil {
  3903. return err
  3904. }
  3905. for _, node := range nodes {
  3906. res := node.Description.Resources
  3907. }
  3908. return nil
  3909. }
  3910. */
  3911. /*
  3912. func numberOfReplicas(ctx context.Context, queueLength int, config shuffle.ScalingConfig) (int, int) {
  3913. queueScaleFactor := queueScaleFactor(queueLength, config)
  3914. numReplicas := int(float64(queueLength) * queueScaleFactor)
  3915. serviceName := "shuffle-workers"
  3916. nodes, err := dockercli.NodeList(ctx, types.NodeListOptions{})
  3917. if err != nil {
  3918. log.Printf("[ERROR] Cannot find any nodes in the swarm network")
  3919. }
  3920. filterArgs := filters.NewArgs()
  3921. filterArgs.Add("service", serviceName)
  3922. filterArgs.Add("desired-state", "running")
  3923. tasks, err := dockercli.TaskList(context.Background(), types.TaskListOptions{
  3924. Filters: filterArgs,
  3925. })
  3926. if err != nil {
  3927. log.Fatalf("[WARNING] Failed to list tasks for service %s: %s", serviceName, err)
  3928. }
  3929. runningReplicas := len(tasks)
  3930. if numReplicas > runningReplicas*len(nodes) {
  3931. maxIncrease := config.MaxScaleUpStep
  3932. if numReplicas > (runningReplicas*len(nodes) + maxIncrease) {
  3933. numReplicas = runningReplicas + maxIncrease
  3934. }
  3935. }
  3936. if numReplicas < config.MinReplicas {
  3937. numReplicas = config.MinReplicas
  3938. }
  3939. if numReplicas > config.MaxReplicas {
  3940. numReplicas = config.MaxReplicas
  3941. }
  3942. return numReplicas, runningReplicas
  3943. }
  3944. */
  3945. // TODO: Currently we use number of request made for the worker to run a execution as it is much
  3946. // easier to track in a window time frame. But this could be useful.
  3947. func collectMetrics(ctx context.Context, dockerClient *dockerclient.Client) (int, error) {
  3948. client := shuffle.GetExternalClient(baseUrl)
  3949. fullUrl := fmt.Sprintf("%s/api/v1/workflows/queue", baseUrl)
  3950. req, err := http.NewRequest("GET", fullUrl, nil)
  3951. if err != nil {
  3952. log.Printf("[ERROR] Failed to send a request to %s: %s", fullUrl, err)
  3953. return 0, err
  3954. }
  3955. req.Header.Add("Content-Type", "application/json")
  3956. req.Header.Add("Org-Id", environment)
  3957. if len(auth) > 0 {
  3958. req.Header.Add("Authorization", auth)
  3959. }
  3960. if len(org) > 0 {
  3961. req.Header.Add("Org", org)
  3962. }
  3963. if len(orborusLabel) > 0 {
  3964. log.Printf("[DEBUG] Sending with Label '%s'", orborusLabel)
  3965. req.Header.Add("X-Orborus-Label", orborusLabel)
  3966. }
  3967. if swarmConfig != "run" && swarmConfig != "swarm" {
  3968. req.Header.Add("X-Orborus-Runmode", "Default")
  3969. } else {
  3970. req.Header.Add("X-Orborus-Runmode", "Docker Swarm")
  3971. }
  3972. resp, err := client.Do(req)
  3973. if err != nil {
  3974. return 0, err
  3975. }
  3976. var executionRequests shuffle.ExecutionRequestWrapper
  3977. body, err := ioutil.ReadAll(resp.Body)
  3978. json.Unmarshal(body, &executionRequests)
  3979. return len(executionRequests.Data), nil
  3980. }
  3981. func setBackendToSwarmNetwork(ctx context.Context) error {
  3982. containerId := ""
  3983. filterArgs := filters.NewArgs()
  3984. filterArgs.Add("name", "shuffle-backend")
  3985. containers, err := dockercli.ContainerList(ctx, container.ListOptions{
  3986. All: true,
  3987. Filters: filterArgs,
  3988. })
  3989. if err != nil {
  3990. return err
  3991. }
  3992. if len(containers) == 0 {
  3993. return errors.New("No containers found with name shuffle-backend")
  3994. }
  3995. containerId = containers[0].ID
  3996. networkName := "shuffle_swarm_executions"
  3997. err = dockercli.NetworkConnect(ctx, networkName, containerId, nil)
  3998. if err != nil {
  3999. log.Printf("[ERROR] Error connecting backend container to network: %s", err)
  4000. }
  4001. return nil
  4002. }