package example;
import java.sql.*;
import java.util.concurrent.TimeUnit;
import static example.Database.*;
public final class Replication {
    private Replication(){}
    static boolean replayed(Connection replica,String lsn)throws SQLException{
        try(var s=replica.prepareStatement("SELECT COALESCE(pg_last_wal_replay_lsn() >= CAST(? AS pg_lsn),false)")){
            s.setString(1,lsn);s.setQueryTimeout(5);try(var r=s.executeQuery()){r.next();return r.getBoolean(1);}
        }
    }
    static void awaitReplay(Connection replica,String lsn)throws Exception{
        long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(10);
        while(!replayed(replica,lsn)){check(System.nanoTime()<deadline,"replica catch-up timeout");Thread.sleep(25);}
    }
    public static void run()throws Exception{
        String url=System.getenv("LAB_REPLICA_URL"), password=System.getenv("LAB_DB_PASSWORD");
        if(url==null||password==null)throw new IllegalStateException("Set LAB_REPLICA_URL and LAB_DB_PASSWORD; start Compose replica profile");
        try(var primary=app(a());var replica=app(url);var admin=connect(url,"lab_owner",password)){
            check("f".equals(scalar(primary,"SELECT pg_is_in_recovery()")),"primary identity");
            check("t".equals(scalar(replica,"SELECT pg_is_in_recovery()")),"replica identity");
            sql(primary,"DELETE FROM tenant_order WHERE id=900");
            awaitReplay(replica,scalar(primary,"SELECT pg_current_wal_lsn()::text"));
            try{
                scalar(admin,"SELECT pg_wal_replay_pause()");
                long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(5);
                while(!"paused".equals(scalar(admin,"SELECT pg_get_wal_replay_pause_state()"))){
                    check(System.nanoTime()<deadline,"pause timeout");Thread.sleep(25);
                }
                sql(primary,"INSERT INTO tenant_order VALUES(900,7,'replica-marker','replica-marker')");
                String target=scalar(primary,"SELECT pg_current_wal_lsn()::text");
                check("1".equals(scalar(primary,"SELECT count(*) FROM tenant_order WHERE id=900")),"primary write");
                check("0".equals(scalar(replica,"SELECT count(*) FROM tenant_order WHERE id=900")),"paused replica old");
                scalar(admin,"SELECT pg_wal_replay_resume()");awaitReplay(replica,target);
                check("1".equals(scalar(replica,"SELECT count(*) FROM tenant_order WHERE id=900")),"replica caught up");
                System.out.println("realStandby=true pausedReplicaOld=true primaryCommitted=true replayLsnCaughtUp=true replicaReadNew=true");
            }finally{
                scalar(admin,"SELECT pg_wal_replay_resume()");
                sql(primary,"DELETE FROM tenant_order WHERE id=900");
                awaitReplay(replica,scalar(primary,"SELECT pg_current_wal_lsn()::text"));
            }
        }
    }
}
