Skip to content

Commit b1517ef

Browse files
Merge pull request #21 from Sreejay-Reddy/feature/CLI-inspect
feat: add inspect API, sentinel CLI, and decouple reconciliation from once()
2 parents c766420 + d72c054 commit b1517ef

17 files changed

Lines changed: 342 additions & 219 deletions

sentinel-py/pyproject.toml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,12 @@ classifiers = [
3535

3636
license = { file = "LICENSE" }
3737

38+
[project.scripts]
39+
sen = "sentinel.cli:main"
40+
3841
[project.optional-dependencies]
3942
django = ["django>=4.2"]
43+
cli = ["python-dotenv"]
4044

4145
[build-system]
4246
requires = ["setuptools>=61.0"]

sentinel-py/sentinel/async_core.py

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import json
22
from .utils import get_owner_id, row_to_dict
3-
from .result import AcquireResult, OperationResult
3+
from .result import AcquireResult, OperationResult, InspectResult
44

55
async def acquire(conn, key, *, owner_id=None, ttl_ms=10000, hard_ttl_ms = None):
66

@@ -182,4 +182,30 @@ async def expire_lease(conn, key, *, owner_id, fencing_token):
182182
await conn.commit()
183183
return OperationResult(success)
184184

185+
async def inspect(conn, key):
186+
async with conn.cursor() as cur:
187+
await cur.execute("""
188+
SELECT owner_id, lease_expires_at, fencing_token, status,
189+
lease_expires_at > NOW() AS lease_alive, lease_updated_at,
190+
hard_expires_at, execution_result
191+
FROM sentinel_leases
192+
WHERE key = %s
193+
""", (key,))
194+
195+
result = await cur.fetchone()
185196

197+
if result is None:
198+
return None
199+
200+
row = row_to_dict(cur, result)
201+
202+
return InspectResult(
203+
key = key,
204+
owner_id = row["owner_id"],
205+
fencing_token = row["fencing_token"],
206+
status = row["status"],
207+
lease_alive = row["lease_alive"],
208+
lease_expires_at = row["lease_expires_at"],
209+
lease_updated_at = row["lease_updated_at"],
210+
hard_expires_at = row["hard_expires_at"],
211+
execution_result = row["execution_result"])

sentinel-py/sentinel/async_once.py

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,6 @@
1313

1414
from .heartbeat_config import get_manager
1515
from .logging import logger
16-
from .async_reconcilliation import AsyncReconcile
1716
from .result import OnceResult
1817

1918

@@ -36,7 +35,6 @@ def __init__(
3635
self.kwargs = kwargs or {}
3736
self.owns_connection = owns_connection
3837

39-
self.reconcile = AsyncReconcile(get_conn)
4038
self._task = None
4139

4240
async def run(self):
@@ -76,8 +74,7 @@ async def run(self):
7674
success=False,
7775
status="executing",
7876
uncertain=True,
79-
execution_alive=False,
80-
reconcile=self.reconcile
77+
execution_alive=False
8178
)
8279

8380
if acquired.status == "executing" and acquired.lease_alive:
@@ -162,8 +159,7 @@ async def run(self):
162159
status="executing",
163160
execution_alive=False,
164161
uncertain=True,
165-
exception=e,
166-
reconcile=self.reconcile
162+
exception=e
167163
)
168164

169165
# Finalize canonical completion
@@ -184,8 +180,7 @@ async def run(self):
184180
return OnceResult(
185181
success=False,
186182
status=completed.status,
187-
uncertain=True,
188-
reconcile=self.reconcile
183+
uncertain=True
189184
)
190185

191186
return OnceResult(

sentinel-py/sentinel/async_reconcilliation.py

Lines changed: 6 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ def __init__(self, get_conn, namespace=None):
66
self.get_conn = get_conn
77
self.namespace = namespace
88

9-
async def reconcile(self, key, *, owner_id, fencing_token):
9+
async def reconcile(self, key):
1010
conn = await self.get_conn()
1111

1212
try:
@@ -17,12 +17,10 @@ async def reconcile(self, key, *, owner_id, fencing_token):
1717
status = 'reconciling',
1818
lease_updated_at = NOW()
1919
WHERE key = %s
20-
AND owner_id = %s
21-
AND fencing_token = %s
2220
AND status = 'executing'
2321
AND lease_expires_at < NOW()
2422
RETURNING 1;
25-
""", (key, owner_id, fencing_token))
23+
""", (key,))
2624

2725
success = await cur.fetchone() is not None
2826

@@ -33,14 +31,7 @@ async def reconcile(self, key, *, owner_id, fencing_token):
3331
finally:
3432
await conn.close()
3533

36-
async def force_complete(
37-
self,
38-
key,
39-
*,
40-
owner_id,
41-
fencing_token,
42-
execution_result
43-
):
34+
async def force_complete(self, key, execution_result):
4435
conn = await self.get_conn()
4536

4637
try:
@@ -52,16 +43,9 @@ async def force_complete(
5243
execution_result = %s,
5344
lease_updated_at = NOW()
5445
WHERE key = %s
55-
AND owner_id = %s
56-
AND fencing_token = %s
5746
AND status = 'reconciling'
5847
RETURNING 1;
59-
""", (
60-
execution_result,
61-
key,
62-
owner_id,
63-
fencing_token
64-
))
48+
""", (execution_result, key))
6549

6650
success = await cur.fetchone() is not None
6751

@@ -72,13 +56,7 @@ async def force_complete(
7256
finally:
7357
await conn.close()
7458

75-
async def reset(
76-
self,
77-
key,
78-
*,
79-
owner_id,
80-
fencing_token
81-
):
59+
async def reset(self, key):
8260
conn = await self.get_conn()
8361

8462
try:
@@ -89,15 +67,9 @@ async def reset(
8967
status = 'claimed',
9068
lease_updated_at = NOW()
9169
WHERE key = %s
92-
AND owner_id = %s
93-
AND fencing_token = %s
9470
AND status = 'reconciling'
9571
RETURNING 1;
96-
""", (
97-
key,
98-
owner_id,
99-
fencing_token
100-
))
72+
""", (key,))
10173

10274
success = await cur.fetchone() is not None
10375

sentinel-py/sentinel/async_sentinel.py

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from .lease import Lease
1+
from .async_reconcilliation import AsyncReconcile
22
from .async_once import AsyncOnce
33
from .heartbeat_config import get_manager
44

@@ -16,6 +16,8 @@ def __init__(self, get_conn = None, default_ttl_ms=3000, namespace=None, owns_co
1616
raise ValueError(
1717
"No database connection provider found."
1818
)
19+
20+
self.reconcile = AsyncReconcile(self._conn, namespace=self.namespace)
1921

2022
self.manager = get_manager(
2123
self._conn,
@@ -35,13 +37,6 @@ def _hard_ttl(self, ttl_ms, hard_ttl_ms):
3537

3638
def _key(self, key):
3739
return f"{self.namespace}:{key}" if self.namespace else key
38-
39-
# def lease(self, key, ttl_ms=None, hard_ttl_ms=None):
40-
# key = self._key(key)
41-
# ttl = self._ttl(ttl_ms)
42-
# hard_ttl = self._hard_ttl(ttl,hard_ttl_ms)
43-
44-
# return Lease(None, key, ttl, hard_ttl, self._conn)
4540

4641
async def once(self, key, fn, ttl_ms=None, hard_ttl_ms=None, kwargs=None):
4742
key = self._key(key)

sentinel-py/sentinel/cli.py

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
import os
2+
import argparse
3+
import psycopg
4+
5+
try:
6+
from dotenv import load_dotenv
7+
load_dotenv()
8+
except ImportError:
9+
pass
10+
11+
from .core import inspect
12+
13+
def get_conn():
14+
db_url = os.environ.get("DATABASE_URL")
15+
if not db_url:
16+
raise EnvironmentError(
17+
"No DATABASE_URL found. Set it in your environment or .env file.\n"
18+
"Example: DATABASE_URL=postgresql://user:pass@localhost/mydb"
19+
)
20+
return psycopg.connect(db_url)
21+
22+
def main():
23+
parser = argparse.ArgumentParser(prog="sen", description="Sentinel CLI")
24+
subparsers = parser.add_subparsers(dest="command")
25+
26+
inspect_parser = subparsers.add_parser("inspect", help="Inspect a lease by key")
27+
inspect_parser.add_argument("key", type=str)
28+
29+
args = parser.parse_args()
30+
31+
if args.command == "inspect":
32+
conn = get_conn()
33+
result = inspect(conn, args.key)
34+
35+
if result is None:
36+
print(f"No lease found for key: {args.key}")
37+
else:
38+
print(f"status: {result.status}")
39+
print(f"lease_alive: {result.lease_alive}")
40+
print(f"owner_id: {result.owner_id}")
41+
print(f"fencing_token: {result.fencing_token}")
42+
print(f"lease_expires_at: {result.lease_expires_at}")
43+
print(f"lease_updated_at: {result.lease_updated_at}")
44+
print(f"hard_expires_at: {result.hard_expires_at}")
45+
print(f"execution_result: {result.execution_result}")
46+
else:
47+
parser.print_help()
48+
49+
if __name__ == "__main__":
50+
main()

sentinel-py/sentinel/core.py

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import json
22
from .utils import get_owner_id, row_to_dict
3-
from .result import AcquireResult, OperationResult
3+
from .result import AcquireResult, OperationResult, InspectResult
44

55
def acquire(conn, key, *, owner_id=None, ttl_ms=10000, hard_ttl_ms = None):
66

@@ -182,4 +182,30 @@ def expire_lease(conn, key, *, owner_id, fencing_token):
182182
conn.commit()
183183
return OperationResult(success)
184184

185+
def inspect(conn, key):
186+
with conn.cursor() as cur:
187+
cur.execute("""
188+
SELECT owner_id, lease_expires_at, fencing_token, status,
189+
lease_expires_at > NOW() AS lease_alive, lease_updated_at,
190+
hard_expires_at, execution_result
191+
FROM sentinel_leases
192+
WHERE key = %s
193+
""", (key,))
194+
195+
result = cur.fetchone()
185196

197+
if result is None:
198+
return None
199+
200+
row = row_to_dict(cur, result)
201+
202+
return InspectResult(
203+
key = key,
204+
owner_id = row["owner_id"],
205+
fencing_token = row["fencing_token"],
206+
status = row["status"],
207+
lease_alive = row["lease_alive"],
208+
lease_expires_at = row["lease_expires_at"],
209+
lease_updated_at = row["lease_updated_at"],
210+
hard_expires_at = row["hard_expires_at"],
211+
execution_result = row["execution_result"])

sentinel-py/sentinel/once.py

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,6 @@
1313

1414
from .heartbeat_config import get_manager
1515
from .logging import logger
16-
from .reconciliation import Reconcile
1716
from .result import OnceResult
1817

1918

@@ -36,7 +35,6 @@ def __init__(
3635
self.kwargs = kwargs or {}
3736
self.owns_connection = owns_connection
3837

39-
self.reconcile = Reconcile(get_conn)
4038
self._task = None
4139

4240
def run(self):
@@ -76,8 +74,7 @@ def run(self):
7674
success=False,
7775
status="executing",
7876
uncertain=True,
79-
execution_alive=False,
80-
reconcile=self.reconcile
77+
execution_alive=False
8178
)
8279

8380
if acquired.status == "executing" and acquired.lease_alive:
@@ -162,8 +159,7 @@ def run(self):
162159
status="executing",
163160
execution_alive=False,
164161
uncertain=True,
165-
exception=e,
166-
reconcile=self.reconcile
162+
exception=e
167163
)
168164

169165
# Finalize canonical completion
@@ -184,8 +180,7 @@ def run(self):
184180
return OnceResult(
185181
success=False,
186182
status=completed.status,
187-
uncertain=True,
188-
reconcile=self.reconcile
183+
uncertain=True
189184
)
190185

191186
return OnceResult(

0 commit comments

Comments
 (0)