MCPcopy Create free account
hub / github.com/cloudwan/gohan / publishEventWithOptions

Function publishEventWithOptions

cli/migrate.go:228–282  ·  view source on GitHub ↗
(envName string, modifiedSchemas []string, eventName string, syncETCDEvent bool, eventTimeout time.Duration, db *server.DbSyncWrapper, manager *schema.Manager, envManager *extension.Manager, sync sync.Sync, ident middleware.IdentityService)

Source from the content-addressed store, hash-verified

226}
227
228func publishEventWithOptions(envName string, modifiedSchemas []string, eventName string, syncETCDEvent bool, eventTimeout time.Duration, db *server.DbSyncWrapper, manager *schema.Manager, envManager *extension.Manager, sync sync.Sync, ident middleware.IdentityService) {
229 deadline := time.Now().Add(eventTimeout)
230
231 for _, s := range manager.Schemas() {
232 if !util.ContainsString(modifiedSchemas, s.ID) {
233 continue
234 }
235
236 pluralURL := s.GetPluralURL()
237
238 if _, ok := envManager.GetEnvironment(s.ID); !ok {
239 now := time.Now()
240 left := deadline.Sub(now)
241 if now.After(deadline) {
242 log.Fatalf("Timeout after '%s' secs while publishing event to schemas", eventTimeout.Seconds())
243 }
244
245 envOtto := otto.NewEnvironment(envName, db, ident, sync)
246 envOtto.SetEventTimeLimit(eventName, left)
247
248 envGoplugin := goplugin.NewEnvironment(envName, nil, nil)
249 envGoplugin.SetDatabase(db)
250 envGoplugin.SetSync(sync)
251
252 env := extension.NewEnvironment([]extension.Environment{envOtto, envGoplugin})
253
254 log.Info("Loading environment for %s schema with URL: %s", s.ID, pluralURL)
255
256 if err := env.LoadExtensionsForPath(manager.Extensions, manager.TimeLimit, manager.TimeLimits, pluralURL); err != nil {
257 log.Fatal(fmt.Sprintf("[%s] %v", pluralURL, err))
258 }
259
260 envManager.RegisterEnvironment(s.ID, env)
261 }
262
263 env, _ := envManager.GetEnvironment(s.ID)
264
265 eventContext := map[string]interface{}{}
266 eventContext["schema"] = s
267 eventContext["schema_id"] = s.ID
268 eventContext["sync"] = sync
269 eventContext["db"] = db
270 eventContext["identity_service"] = ident
271
272 if err := env.HandleEvent(eventName, eventContext); err != nil {
273 log.Fatalf("Failed to handle event '%s': %s", eventName, err)
274 }
275 }
276
277 if syncETCDEvent {
278 if _, err := server.NewSyncWriter(sync, db).Sync(); err != nil {
279 log.Fatalf("Failed to synchronize post-migration events, err: %s", err)
280 }
281 }
282}
283
284func publishEvent(envName string, modifiedSchemas []string, eventName string, syncETCDEvent bool, eventTimeout time.Duration) error {
285 config := util.GetConfig()

Callers 1

publishEventFunction · 0.85

Calls 14

SetEventTimeLimitMethod · 0.95
SetDatabaseMethod · 0.95
SetSyncMethod · 0.95
LoadExtensionsForPathMethod · 0.95
HandleEventMethod · 0.95
GetPluralURLMethod · 0.80
GetEnvironmentMethod · 0.80
FatalfMethod · 0.80
FatalMethod · 0.80
RegisterEnvironmentMethod · 0.80
SchemasMethod · 0.65
InfoMethod · 0.65

Tested by

no test coverage detected