Repository navigation
Expand file tree
/
Copy pathsubclass_demo.py
More file actions
67 lines (57 loc) · 2.42 KB
/
Copy pathsubclass_demo.py
File metadata and controls
67 lines (57 loc) · 2.42 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
from jms_pam import Client
from jms_pam_config import client_options, instance_id
def apply_account(account):
# Validate a new connection, switch the pool, then release old connections.
# Never log account.secret or authentication headers.
raise NotImplementedError("Implement the application's account update first")
def restart_application():
raise NotImplementedError("Implement application restart and health check first")
class ApplicationClient(Client):
def apply_event(self, event):
if event.get("event") == "application.restart.requested":
restart_application()
return
account = self.get_account(
account_id=event["account_id"], allow_local_fallback=False,
)
if account.revision != event["account_revision"]:
raise ValueError("The event account version is superseded")
apply_account(account)
def handle_event(self, event):
event_id = event["event_id"]
if event.get("command_id"):
claim = self.report_application_command_result(
command_id=event["command_id"], status="running",
)
if not claim.accepted:
return
try:
self.apply_event(event)
except Exception:
self.confirm_event(
event_id=event_id, status="failed", error_code="application_failed",
)
raise
# A fetch or delivery receipt never confirms application success.
self.confirm_event(event_id=event_id)
def run(self):
for event in self.watch_credential_events():
if event.get("command_id"):
updates = [event]
elif event.get("event") == "snapshot":
updates = event.get("credentials", [])
# Release application connections absent from the new scope.
elif event.get("event") == "credential.updated":
updates = [event]
elif event.get("event") == "credential.revoked":
# Release connections for event["account_id"].
continue
else:
continue
for update in updates:
try:
self.handle_event(update)
except Exception as error:
self.on_event_error(error, update)
with ApplicationClient(instance_id=instance_id, **client_options) as client:
client.run()