(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)
| 226 | } |
| 227 | |
| 228 | func 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 | |
| 284 | func publishEvent(envName string, modifiedSchemas []string, eventName string, syncETCDEvent bool, eventTimeout time.Duration) error { |
| 285 | config := util.GetConfig() |
no test coverage detected