"""Bounded polling example. Run as a script; see README.md for setup.""" import argparse import asyncio import fcntl import json import os from pathlib import Path import sqlite3 import tempfile import uuid import httpx def checkpoint(path,value): path=Path(path);path.parent.mkdir(parents=True,exist_ok=True) fd,temp=tempfile.mkstemp(prefix=path.name+'.',dir=path.parent) try: with os.fdopen(fd,'w') as stream: json.dump(value,stream);stream.flush();os.fsync(stream.fileno()) os.replace(temp,path) directory=os.open(path.parent,os.O_RDONLY) try:os.fsync(directory) finally:os.close(directory) finally: if os.path.exists(temp):os.unlink(temp) class LocalSink: """Idempotent demo effect: retain event IDs/status only, not Reddit content.""" def __init__(self,path):self.path=path async def __call__(self,items): with sqlite3.connect(self.path) as db: db.execute('CREATE TABLE IF NOT EXISTS seen(id TEXT PRIMARY KEY,status TEXT NOT NULL)') for item in items: db.execute('INSERT INTO seen VALUES (?,?) ON CONFLICT(id) DO UPDATE SET status=excluded.status',(item['id'],item['status'])) async def poll(client,*,state_path,payload,handle,max_pages=5): if not 1<=max_pages<=10:raise ValueError('One to ten pages per invocation.') path=Path(state_path);path.parent.mkdir(parents=True,exist_ok=True) # Reject overlapping runs instead of blocking a scheduler indefinitely. with open(str(path)+'.lock','a') as lock: fcntl.flock(lock,fcntl.LOCK_EX|fcntl.LOCK_NB) state=json.loads(path.read_text()) if path.exists() else {'intent':str(uuid.uuid4()),'payload':payload,'base':str(client.base_url)} if state['payload']!=payload or state['base']!=str(client.base_url):raise ValueError('Use a new checkpoint for a different monitor or API.') checkpoint(path,state) # persist intent BEFORE a request that may time out if not state.get('monitor_id'): response=await client.post('/api/v1/monitors',json=payload,headers={'Idempotency-Key':state['intent']}) response.raise_for_status();state['monitor_id']=response.json()['id'];checkpoint(path,state) processed=0 for _ in range(max_pages): params={'limit':25} if state.get('cursor'):params['after']=state['cursor'] response=await client.get(f"/api/v1/monitors/{state['monitor_id']}/updates",params=params) response.raise_for_status();page=response.json() await handle(page['items']) # failures leave cursor unchanged; effects must dedupe on item.id state['cursor']=page['next_cursor'];checkpoint(path,state) processed+=len(page['items']) if not page['has_more']:break return {'monitor_id':state['monitor_id'],'processed':processed,'has_more':page['has_more'], 'latest_check':page['monitor']['latest_check']} async def main(args): key=os.environ.get('MONITOR_API_KEY') if not key:raise SystemExit('Set MONITOR_API_KEY in your environment.') async with httpx.AsyncClient(base_url=args.base,headers={'Authorization':'Bearer '+key},timeout=20,follow_redirects=False) as client: result=await poll(client,state_path=args.state,payload={'name':args.name,'query':args.query,'communities':args.community}, handle=LocalSink(str(args.state)+'.sqlite'),max_pages=args.max_pages) print(json.dumps(result,indent=2)) if __name__=='__main__': parser=argparse.ArgumentParser(description=__doc__) parser.add_argument('--base',default='http://127.0.0.1:8020') parser.add_argument('--state',type=Path,default=Path('output/agent-example/checkpoint.json')) parser.add_argument('--name',default='Tool discussions');parser.add_argument('--query',default='tool') parser.add_argument('--community',action='append',default=None);parser.add_argument('--max-pages',type=int,default=5) args=parser.parse_args();args.community=args.community or ['fixture_webdev'] asyncio.run(main(args))